diff --git a/.gitignore b/.gitignore index dfcf1b0..b1f1f7e 100644 --- a/.gitignore +++ b/.gitignore @@ -4,6 +4,7 @@ __pycache__/ *.pyo .venv/ venv/ +.pytest_cache/ # Entorno / secretos .env @@ -15,6 +16,7 @@ estado/*.log estado/.watcher_last_id estado/.netatmo_last_alert.json estado/.netatmo_token.json +estado/monitor/ # IDE / OS .vscode/ diff --git a/CHANGELOG_INFRA.md b/CHANGELOG_INFRA.md index 2ca39e0..474ce08 100644 --- a/CHANGELOG_INFRA.md +++ b/CHANGELOG_INFRA.md @@ -5,10 +5,10 @@ ## Pendientes de replicar -- **2026-09-29-migraci-n-de-watchers-al-vps** [servicios] Migración de watchers al VPS (aplicado en: sobremesa) - **2026-09-30-fix-telegram-watcher-telethon-1-44-encoding-utf-8** [watchers] Fix telegram_watcher (telethon 1.44 + encoding UTF-8) (aplicado en: portatil) - **2026-10-02-superguardado-py-protocolo-guardado-total-automatico-portati** [scripts] superguardado.py: protocolo guardado total automatico portatil+sobremesa (aplicado en: portatil) - **2026-10-03-superwhisper-enter-ahk-f13-boton-auxiliar-logitech-tarea-pro** [scripts] SuperWhisper Enter: AHK F13 boton auxiliar Logitech + tarea programada (aplicado en: portatil) +- **2026-10-03-monitor-de-estado-telegram-correo-jira-en-segundo-plano** [servicios] Monitor de estado (Telegram + correo + Jira) en segundo plano (aplicado en: portatil) ## Historial @@ -20,6 +20,14 @@ - **Aplicado en**: portatil - **Detalle**: Script superwhisper_enter.ahk: F13 (boton auxiliar Logitech asignado en Logi Options+) espera cambio portapapeles (Super Whisper pega) y manda Enter + devuelve foco a OpenCode. Tarea programada SuperWhisper Enter al inicio de sesion. Rueda central = toggle Super Whisper (inicio/parar). F13 = Enter automatico. +### 2026-10-03 · Monitor de estado (Telegram + correo + Jira) en segundo plano +- **id**: `2026-10-03-monitor-de-estado-telegram-correo-jira-en-segundo-plano` +- **Equipo**: portatil +- **Categoria**: servicios +- **Transferible**: si +- **Aplicado en**: portatil +- **Detalle**: Proceso Python (scripts/monitor_estado.py) que cada 30 min recoge Telegram (5 fuentes), correo jminguez@prolongo.es y Jira DS (nuevos + sin asignar), archiva por fuente, resume con IA (OmniRoute) y avisa por voz local. Arranque por tarea programada 'Monitor Estado' + plugin global monitor-estado-bootstrap.ts; watchdog idempotente; sesion Telegram dedicada (sessions/monitor, login por QR); tests pytest en tests/. Requiere en el otro equipo: pip install pytest qrcode, login Telegram QR y re-ejecutar instalar_global.ps1. + ### 2026-10-02 · superguardado.py: protocolo guardado total automatico portatil+sobremesa - **id**: `2026-10-02-superguardado-py-protocolo-guardado-total-automatico-portati` - **Equipo**: portatil @@ -41,7 +49,7 @@ - **Equipo**: sobremesa - **Categoria**: servicios - **Transferible**: si -- **Aplicado en**: sobremesa +- **Aplicado en**: sobremesa, portatil - **Detalle**: Telegram Watcher y Shutdown Listener migrados al VPS. Recordatorios Msgbox eliminado. Ver INFRAESTRUCTURA/vps/MIGRACION_SERVICIOS_VPS.md para instrucciones. ### 2026-09-20 · OmniRoute arranque fiable (tarea programada + vigilante) diff --git a/changelog_infra.json b/changelog_infra.json index 8fcd175..53642c1 100644 --- a/changelog_infra.json +++ b/changelog_infra.json @@ -185,7 +185,8 @@ "titulo": "Migración de watchers al VPS", "transferible": true, "aplicado_en": [ - "sobremesa" + "sobremesa", + "portatil" ], "detalle": "Telegram Watcher y Shutdown Listener migrados al VPS. Recordatorios Msgbox eliminado. Ver INFRAESTRUCTURA/vps/MIGRACION_SERVICIOS_VPS.md para instrucciones." }, @@ -224,6 +225,18 @@ "portatil" ], "detalle": "Script superwhisper_enter.ahk: F13 (boton auxiliar Logitech asignado en Logi Options+) espera cambio portapapeles (Super Whisper pega) y manda Enter + devuelve foco a OpenCode. Tarea programada SuperWhisper Enter al inicio de sesion. Rueda central = toggle Super Whisper (inicio/parar). F13 = Enter automatico." + }, + { + "id": "2026-10-03-monitor-de-estado-telegram-correo-jira-en-segundo-plano", + "fecha": "2026-10-03", + "equipo": "portatil", + "categoria": "servicios", + "titulo": "Monitor de estado (Telegram + correo + Jira) en segundo plano", + "transferible": true, + "aplicado_en": [ + "portatil" + ], + "detalle": "Proceso Python (scripts/monitor_estado.py) que cada 30 min recoge Telegram (5 fuentes), correo jminguez@prolongo.es y Jira DS (nuevos + sin asignar), archiva por fuente, resume con IA (OmniRoute) y avisa por voz local. Arranque por tarea programada 'Monitor Estado' + plugin global monitor-estado-bootstrap.ts; watchdog idempotente; sesion Telegram dedicada (sessions/monitor, login por QR); tests pytest en tests/. Requiere en el otro equipo: pip install pytest qrcode, login Telegram QR y re-ejecutar instalar_global.ps1." } ] } diff --git a/opencode/monitor-estado-bootstrap.ts b/opencode/monitor-estado-bootstrap.ts new file mode 100644 index 0000000..63507c6 --- /dev/null +++ b/opencode/monitor-estado-bootstrap.ts @@ -0,0 +1,67 @@ +/** + * Monitor Estado Bootstrap — plugin GLOBAL de OpenCode + * + * Objetivo: que al arrancar OpenCode (en CUALQUIER repositorio) se asegure de que el + * Monitor de estado esta corriendo, lanzando su vigilante en modo `-Once` (idempotente: + * el propio monitor tiene un lock, asi que nunca habra dos instancias). + * + * No depende de ningun repositorio: vive en la config global de OpenCode + * (~/.config/opencode/plugins/) y el vigilante canonico esta en el repo INFRAESTRUCTURA. + * + * Instalacion: `scripts/instalar_global.ps1`. + * + * Nota: el vigilante ya corre como tarea programada al iniciar sesion; este plugin es + * la red de seguridad para cuando OpenCode arranca sin que el monitor este vivo. + */ + +import type { Plugin } from "@opencode-ai/plugin" +import { spawn } from "node:child_process" +import { existsSync } from "node:fs" + +// --- Configuracion --------------------------------------------------------- + +// Vigilante canonico (repo INFRAESTRUCTURA). Se puede sobreescribir con MONITOR_ESTADO_WATCHDOG. +const WATCHDOG = + process.env.MONITOR_ESTADO_WATCHDOG ?? + "C:\\Users\\juanm\\Documents\\Github\\INFRAESTRUCTURA\\scripts\\monitor_estado_watchdog.ps1" + +const POWERSHELL = + process.env.SystemRoot + ? `${process.env.SystemRoot}\\System32\\WindowsPowerShell\\v1.0\\powershell.exe` + : "C:\\Windows\\System32\\WindowsPowerShell\\v1.0\\powershell.exe" + +// Evita lanzar el vigilante mas de una vez por proceso de OpenCode. +let launched = false + +function ensureMonitor(): void { + if (launched) return + launched = true + try { + if (!existsSync(WATCHDOG)) return + const child = spawn( + POWERSHELL, + [ + "-NoProfile", + "-NonInteractive", + "-ExecutionPolicy", + "Bypass", + "-WindowStyle", + "Hidden", + "-File", + WATCHDOG, + "-Once", + ], + { detached: true, stdio: "ignore", windowsHide: true }, + ) + child.unref() + } catch { + // Silencioso: el arranque de OpenCode no debe romperse por el monitor. + } +} + +// --- Plugin ---------------------------------------------------------------- + +export const MonitorEstadoBootstrap: Plugin = async () => { + ensureMonitor() + return {} +} diff --git a/scripts/instalar_global.ps1 b/scripts/instalar_global.ps1 index c5289b7..f1609f6 100644 --- a/scripts/instalar_global.ps1 +++ b/scripts/instalar_global.ps1 @@ -31,6 +31,12 @@ $PluginDir = Join-Path $env:USERPROFILE '.config\opencode\plugins' $PluginDst = Join-Path $PluginDir 'omniroute-bootstrap.ts' $UserId = "$env:USERDOMAIN\$env:USERNAME" +# Monitor de estado (Telegram + correo + Jira) +$TaskMonitor = 'Monitor Estado' +$WatchdogMonitor = Join-Path $RepoRoot 'scripts\monitor_estado_watchdog.ps1' +$PluginSrcMon = Join-Path $RepoRoot 'opencode\monitor-estado-bootstrap.ts' +$PluginDstMon = Join-Path $PluginDir 'monitor-estado-bootstrap.ts' + Write-Host '== Instalando infraestructura global (OmniRoute + OpenCode) ==' -ForegroundColor Cyan # 1) Plugin global de OpenCode --------------------------------------------- @@ -38,6 +44,9 @@ if (-not (Test-Path $PluginDir)) { New-Item -ItemType Directory -Force -Path $Pl Copy-Item -LiteralPath $PluginSrc -Destination $PluginDst -Force Write-Host " - Plugin global instalado en $PluginDst" -ForegroundColor Green +Copy-Item -LiteralPath $PluginSrcMon -Destination $PluginDstMon -Force +Write-Host " - Plugin global instalado en $PluginDstMon" -ForegroundColor Green + # 2) Tarea programada ------------------------------------------------------ if (-not (Test-Path $Watchdog)) { throw "No se encuentra el vigilante: $Watchdog" } @@ -76,10 +85,43 @@ Register-ScheduledTask ` Write-Host " - Tarea '$TaskName' registrada -> $Watchdog" -ForegroundColor Green +# 2b) Tarea programada del Monitor de estado -------------------------------- +if (-not (Test-Path $WatchdogMonitor)) { throw "No se encuentra el vigilante: $WatchdogMonitor" } + +if (Get-ScheduledTask -TaskName $TaskMonitor -ErrorAction SilentlyContinue) { + Write-Host ' - Eliminando tarea anterior del monitor...' -ForegroundColor Yellow + Unregister-ScheduledTask -TaskName $TaskMonitor -Confirm:$false +} + +$ActionMon = New-ScheduledTaskAction ` + -Execute $PwshExe ` + -Argument "-NoProfile -NonInteractive -ExecutionPolicy Bypass -WindowStyle Hidden -File `"$WatchdogMonitor`" -Once" ` + -WorkingDirectory $env:USERPROFILE + +$SettingsMon = New-ScheduledTaskSettingsSet ` + -AllowStartIfOnBatteries ` + -DontStopIfGoingOnBatteries ` + -StartWhenAvailable ` + -MultipleInstances IgnoreNew +$SettingsMon.ExecutionTimeLimit = 'PT0S' +$SettingsMon.Hidden = $true + +Register-ScheduledTask ` + -TaskName $TaskMonitor ` + -Action $ActionMon ` + -Trigger $Trigger ` + -Settings $SettingsMon ` + -Principal $Principal ` + -Description 'Arranca el Monitor de estado (Telegram + correo + Jira) al iniciar sesion.' ` + -Force | Out-Null + +Write-Host " - Tarea '$TaskMonitor' registrada -> $WatchdogMonitor" -ForegroundColor Green + # 3) Arrancar -------------------------------------------------------------- if (-not $NoStart) { Start-ScheduledTask -TaskName $TaskName - Write-Host ' - Vigilante arrancado.' -ForegroundColor Yellow + Start-ScheduledTask -TaskName $TaskMonitor + Write-Host ' - Vigilantes arrancados.' -ForegroundColor Yellow } # 4) Watchers de infraestructura (Telegram + recordatorios) ----------------- diff --git a/scripts/monitor_estado.py b/scripts/monitor_estado.py new file mode 100644 index 0000000..d44f036 --- /dev/null +++ b/scripts/monitor_estado.py @@ -0,0 +1,215 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +"""monitor_estado.py — Punto de entrada del Monitor de estado. + +Modos: + ciclo Ejecuta un ciclo completo y termina (a demanda). + daemon Bucle en segundo plano cada `intervalo_minutos` (defecto 30). + estado Muestra el ultimo resumen y las rutas de los ficheros. + login-telegram Autoriza la sesion dedicada de Telegram (una vez). + +El nucleo vive en el paquete `monitor_estado/`. +""" + +from __future__ import annotations + +import argparse +import sys +import time +from datetime import datetime +from pathlib import Path + +# Permite ejecutar el script desde cualquier directorio. +_SCRIPTS_DIR = Path(__file__).resolve().parent +if str(_SCRIPTS_DIR) not in sys.path: + sys.path.insert(0, str(_SCRIPTS_DIR)) + +import monitor_estado.config as config # noqa: E402 +from monitor_estado import __version__, ciclo, lock # noqa: E402 + +# Salida UTF-8 en consola Windows (acentos y simbolos). +for _stream in (sys.stdout, sys.stderr): + try: + _stream.reconfigure(encoding="utf-8", errors="replace") + except Exception: + pass + + +class _NoJira: + """Cliente Jira vacio (para el modo --solo que desactiva esta fuente).""" + + def search(self, jql, fields=None, max_results=50): + return {"issues": []} + + +def _deps_solo(solo: str | None) -> dict: + """Dependencias que desactivan las fuentes no seleccionadas con --solo.""" + deps: dict = {} + if solo and solo != "telegram": + deps["telegram_recoger"] = lambda cfg, hoy: ({}, []) + if solo and solo != "correo": + deps["correo_recoger"] = lambda horas=24, limite=50: [] + if solo and solo != "jira": + deps["jira_cliente"] = _NoJira() + return deps + + +def cmd_ciclo(args: argparse.Namespace) -> int: + """Ejecuta un ciclo completo (o de una sola fuente con --solo).""" + cfg = config.cargar() + if not lock.adquirir(): + print("[monitor_estado] Ya hay una instancia en marcha; no lanzo otra.") + return 0 + try: + ciclo.ejecutar(cfg, deps=_deps_solo(getattr(args, "solo", None))) + except Exception as e: + print(f"[monitor_estado] Error en el ciclo: {e}", file=sys.stderr) + return 1 + finally: + lock.liberar() + print(ciclo.texto_estado()) + return 0 + + +def cmd_daemon(args: argparse.Namespace) -> int: + """Bucle en segundo plano cada intervalo_minutos (RF-19).""" + if not lock.adquirir(): + print("[monitor_estado] Ya hay una instancia en marcha; salgo.") + return 0 + cfg = config.cargar() + intervalo = max(1, int(cfg.get("intervalo_minutos", 30))) * 60 + print(f"[monitor_estado] daemon iniciado (cada {intervalo // 60} min). Ctrl+C para salir.", flush=True) + try: + while True: + try: + res = ciclo.ejecutar(cfg) + atencion = len(res.get("requiere_atencion", [])) + print( + f"[{datetime.now():%Y-%m-%d %H:%M:%S}] ciclo ok | fuentes={res.get('fuentes_con_novedad')} " + f"| atencion={atencion} | aviso={res.get('aviso_emitido')}", + flush=True, + ) + except Exception as e: # un ciclo fallido no debe matar el daemon + print(f"[{datetime.now():%Y-%m-%d %H:%M:%S}] Error en ciclo: {e}", flush=True, file=sys.stderr) + time.sleep(intervalo) + except KeyboardInterrupt: + pass + finally: + lock.liberar() + return 0 + + +def cmd_estado(args: argparse.Namespace) -> int: + """Muestra el ultimo resumen y la configuracion (RF-18).""" + print(f"[monitor_estado] v{__version__}") + print(ciclo.texto_estado()) + print("\n--- Configuracion ---") + print(config.resumen_config(config.cargar())) + return 0 + + +def cmd_simular(args: argparse.Namespace) -> int: + """Simula entradas ficticias (correo/Telegram) SIN enviar ni crear nada real. + + Ejecuta el pipeline completo en un directorio aislado (`estado/monitor/simulacion/`): + no toca los datos reales, no crea correos y NUNCA envia mensajes a Telegram. + """ + from datetime import date + + import monitor_estado.archivo as archivo + import monitor_estado.vistos as vistos + + base = config.MONITOR_DIR / "simulacion" + archivo.MONITOR_DIR = base + vistos.VISTO_DIR = base / "visto" + hoy = date.today() + texto = getattr(args, "texto", None) or "Confirmado el desarrollo" + asunto = getattr(args, "asunto", None) or texto + + def tg_fake(cfg, h): + return ({"simulacion_telegram": [{ + "id": 9900001, "fecha": datetime.now().isoformat(), + "quien": getattr(args, "de", None) or "DAX", "texto": texto, + }]}, []) + + def correo_fake(horas=24, limite=50): + remitente = getattr(args, "de", None) or "DAX" + return [{ + "id": "sim-mail-1", "fecha": datetime.now().strftime("%a, %d %b %Y %H:%M:%S +0200"), + "de": f"{remitente} <{remitente.lower()}@prolongo.es>", + "asunto": asunto, "snippet": texto, + }] + + res = ciclo.ejecutar( + cfg=config.cargar(), hoy=hoy, + deps={"telegram_recoger": tg_fake, "correo_recoger": correo_fake, "jira_cliente": _NoJira()}, + ) + print("[SIMULACION] No se ha enviado ni creado nada real (solo simulado, aislado en 'simulacion/').") + print(f"Elementos de atencion: {len(res.get('requiere_atencion', []))} | aviso={res.get('aviso_emitido')}") + print(ciclo.texto_estado()) + return 0 + + +def cmd_login_telegram(args: argparse.Namespace) -> int: + """Autoriza la sesion dedicada de Telegram (interactivo, una sola vez).""" + from monitor_estado.fuentes import telegram as tg + + if getattr(args, "qr", False): + tg.login_qr(config.cargar()) + else: + tg.login(config.cargar(), force_sms=getattr(args, "sms", False)) + return 0 + + +def build_parser() -> argparse.ArgumentParser: + """Construye el parser de argumentos del CLI.""" + parser = argparse.ArgumentParser( + prog="monitor_estado", + description="Monitor de estado (Telegram + correo + Jira) con resumen IA y aviso local.", + ) + parser.add_argument("--version", action="version", version=f"monitor_estado {__version__}") + subs = parser.add_subparsers(dest="comando") + + p_ciclo = subs.add_parser("ciclo", help="Ejecuta un ciclo completo y termina.") + p_ciclo.add_argument("--solo", choices=["telegram", "correo", "jira"], help="Recoge solo esa fuente.") + + subs.add_parser("daemon", help="Bucle en segundo plano cada intervalo_minutos.") + + subs.add_parser("estado", help="Muestra el ultimo resumen y las rutas.") + + p_sim = subs.add_parser("simular", help="Simula entradas ficticias sin enviar ni crear nada real.") + p_sim.add_argument("--texto", help="Texto simulado (defecto: 'Confirmado el desarrollo').") + p_sim.add_argument("--asunto", help="Asunto del correo simulado.") + p_sim.add_argument("--de", help="Remitente simulado (defecto: DAX).") + + p_login = subs.add_parser("login-telegram", help="Autoriza la sesion dedicada de Telegram (una vez).") + p_login.add_argument("--sms", action="store_true", help="Forzar que el codigo llegue por SMS.") + p_login.add_argument("--qr", action="store_true", help="Autorizar escaneando un codigo QR.") + return parser + + +def main(argv: list[str] | None = None) -> int: + """Punto de entrada: despacha el subcomando o muestra ayuda + config.""" + parser = build_parser() + args = parser.parse_args(argv) + + if args.comando == "ciclo": + return cmd_ciclo(args) + if args.comando == "daemon": + return cmd_daemon(args) + if args.comando == "estado": + return cmd_estado(args) + if args.comando == "simular": + return cmd_simular(args) + if args.comando == "login-telegram": + return cmd_login_telegram(args) + + # Sin subcomando: ayuda + resumen de configuracion. + parser.print_help() + print("\n--- Configuracion actual ---") + print(config.resumen_config(config.cargar())) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/monitor_estado/README.md b/scripts/monitor_estado/README.md new file mode 100644 index 0000000..b43fff0 --- /dev/null +++ b/scripts/monitor_estado/README.md @@ -0,0 +1,58 @@ +# Monitor de estado — Telegram + correo + Jira + +Proceso en segundo plano que **cada 30 minutos** recoge lo del día en curso, lo archiva +por fuente, genera un **resumen IA** y emite un **aviso por voz local** cuando algo +requiere atención. Nunca envía mensajes a ningún medio (solo lee y avisa localmente). + +Spec y plan: `biblioteca_negocio_prolongo/specs/005-monitor-estado-conversaciones/`. + +## Qué recoge +- **Telegram** (sesión dedicada `~/.telegram-agent/sessions/monitor`): Oscar Rodríguez, Isa, + Daniel García Jiménez, Antonio Romero y el canal "Desarrollo software". +- **Correo** `jminguez@prolongo.es`: últimos 24 h (no marca leído). +- **Jira** proyecto **DS**: incidencias nuevas (desde el ciclo anterior) y tickets sin asignar + no cerrados. Solo lectura. + +## Uso +```powershell +python scripts/monitor_estado.py ciclo # un ciclo y salir +python scripts/monitor_estado.py ciclo --solo correo +python scripts/monitor_estado.py daemon # bucle cada 30 min (segundo plano) +python scripts/monitor_estado.py estado # último resumen + rutas +python scripts/monitor_estado.py login-telegram --qr # autorizar sesión Telegram (una vez) +python scripts/monitor_estado.py simular --de DAX --texto "Confirmado el desarrollo" +``` + +- `simular` inyecta entradas ficticias en un directorio aislado (`estado/monitor/simulacion/`): + no crea correos reales, no toca datos reales y **nunca** envía a Telegram. + +## Arranque y vigilancia +- **Tarea programada `Monitor Estado`** (al iniciar sesión) → `scripts/monitor_estado_watchdog.ps1 -Once`. +- **Plugin global de OpenCode** `~/.config/opencode/plugins/monitor-estado-bootstrap.ts` → relanza el + watchdog al arrancar OpenCode. +- El watchdog lanza `daemon` solo si no corre; el propio monitor tiene un lock (instancia única). +- Instalación: `scripts/instalar_global.ps1`. +```powershell +powershell -ExecutionPolicy Bypass -File scripts/monitor_estado_watchdog.ps1 -Once +powershell -ExecutionPolicy Bypass -File scripts/monitor_estado_watchdog.ps1 -Stop +``` + +## Configuración +`estado/monitor/config.json` (se autogenera desde `config.example.json`): +intervalo, retención (30 días), fuentes de Telegram, correo, Jira, IA (OmniRoute) y aviso de voz. + +## Datos y logs +``` +estado/monitor/ + conversaciones/ correo/ jira/ # ficheros del día por fuente (.md) + visto/ # estado "ya visto" (dedup) + resumen_actual.md/.json # último resumen + ultimo_ciclo.json + monitor_estado.out.log / .err.log # salida del daemon +``` + +## Tests +```powershell +python -m pytest tests/ -q +``` +Requiere: `pip install pytest qrcode` (Telethon/gspread ya usados por otros watchers). diff --git a/scripts/monitor_estado/__init__.py b/scripts/monitor_estado/__init__.py new file mode 100644 index 0000000..0b96532 --- /dev/null +++ b/scripts/monitor_estado/__init__.py @@ -0,0 +1,11 @@ +# -*- coding: utf-8 -*- +"""monitor_estado — Monitor de estado en segundo plano. + +Paquete que recoge cada 30 minutos lo del dia en curso de Telegram (5 fuentes), +correo de empresa y Jira DS, lo archiva por fuente, genera un resumen IA y emite +un aviso por voz local cuando algo requiere atencion. + +Ver spec `specs/005-monitor-estado-conversaciones/` en el repo de negocio. +""" + +__version__ = "0.1.0" diff --git a/scripts/monitor_estado/archivo.py b/scripts/monitor_estado/archivo.py new file mode 100644 index 0000000..bda2192 --- /dev/null +++ b/scripts/monitor_estado/archivo.py @@ -0,0 +1,122 @@ +# -*- coding: utf-8 -*- +"""archivo.py — Archivado diario por fuente y resumen (RF-5, RF-7, RF-9, RF-10, RF-18). + +Escribe un fichero Markdown por fuente y dia (`estado/monitor//.md`). +En cada ciclo el fichero del dia se **reescribe** con el estado actual completo, de modo +que siempre contiene todo lo de hoy y nada duplicado. Los ficheros de mas de +`retencion_dias` se eliminan. Tambien persiste el ultimo resumen. + +El formateo de cada elemento se delega en `render_fn`, para no acoplar este modulo a +la forma concreta de cada fuente. +""" + +from __future__ import annotations + +import json +import re +from datetime import date, datetime, timedelta +from pathlib import Path +from typing import Callable, Iterable + +from . import config + +MONITOR_DIR: Path = config.MONITOR_DIR + +_PATRON_FECHA = re.compile(r"^(\d{4}-\d{2}-\d{2})\.md$") + + +def _fecha_str(fecha) -> str: + """Normaliza una fecha (date o texto ISO) a 'YYYY-MM-DD'.""" + if isinstance(fecha, datetime): + fecha = fecha.date() + if isinstance(fecha, date): + return fecha.isoformat() + return str(fecha) + + +def escribe_dia( + ruta_logica: str, + fecha, + elementos: Iterable, + id_fn: Callable = lambda e: e, + render_fn: Callable = str, + titulo: str | None = None, +) -> Path: + """Reescribe el fichero del dia de una fuente, sin duplicados por id. + + `ruta_logica` es la ruta relativa bajo el directorio del monitor (p. ej. + `telegram/oscar_rodriguez`, `correo`, `jira`). + """ + fecha_str = _fecha_str(fecha) + vistos: set[str] = set() + bloques: list[str] = [] + for e in elementos: + clave = str(id_fn(e)) + if clave in vistos: + continue + vistos.add(clave) + bloques.append(render_fn(e)) + + titulo = titulo or f"{ruta_logica} — {fecha_str}" + cuerpo = "\n\n".join(bloques) + contenido = f"# {titulo}\n\n{cuerpo}\n" if cuerpo else f"# {titulo}\n" + + destino = MONITOR_DIR / ruta_logica / f"{fecha_str}.md" + destino.parent.mkdir(parents=True, exist_ok=True) + destino.write_text(contenido, encoding="utf-8") + return destino + + +def poda(retencion_dias: int = 30, hoy: date | None = None) -> list[Path]: + """Elimina los ficheros diarios con fecha anterior a `hoy - retencion_dias`. + + Devuelve la lista de ficheros borrados. + """ + hoy = hoy or date.today() + limite = hoy - timedelta(days=retencion_dias) + borrados: list[Path] = [] + if not MONITOR_DIR.exists(): + return borrados + for f in MONITOR_DIR.rglob("*.md"): + m = _PATRON_FECHA.match(f.name) + if not m: + continue + try: + fd = date.fromisoformat(m.group(1)) + except ValueError: + continue + if fd < limite: + try: + f.unlink() + borrados.append(f) + except OSError: + pass + return borrados + + +def escribe_resumen(resumen: dict) -> None: + """Persiste el ultimo resumen en `resumen_actual.md` y `resumen_actual.json`.""" + MONITOR_DIR.mkdir(parents=True, exist_ok=True) + md = resumen.get("resumen_md", "") or "" + (MONITOR_DIR / "resumen_actual.md").write_text(md, encoding="utf-8") + with open(MONITOR_DIR / "resumen_actual.json", "w", encoding="utf-8") as f: + json.dump(resumen, f, indent=2, ensure_ascii=False) + + +def escribe_ultimo_ciclo(data: dict) -> None: + """Persiste la marca del ultimo ciclo en `ultimo_ciclo.json`.""" + MONITOR_DIR.mkdir(parents=True, exist_ok=True) + with open(MONITOR_DIR / "ultimo_ciclo.json", "w", encoding="utf-8") as f: + json.dump(data, f, indent=2, ensure_ascii=False) + + +def leer_ultimo_ciclo() -> dict: + """Lee `ultimo_ciclo.json` (o {} si no existe).""" + ruta = MONITOR_DIR / "ultimo_ciclo.json" + if ruta.exists(): + try: + with open(ruta, "r", encoding="utf-8") as f: + return json.load(f) + except Exception: + pass + return {} diff --git a/scripts/monitor_estado/aviso.py b/scripts/monitor_estado/aviso.py new file mode 100644 index 0000000..b601bed --- /dev/null +++ b/scripts/monitor_estado/aviso.py @@ -0,0 +1,53 @@ +# -*- coding: utf-8 -*- +"""aviso.py — Aviso por voz local (RF-14, RF-15). + +Emite un aviso sonoro local (reutiliza `scripts/aviso_sonoro.py` del repo de negocio) +solo cuando el resumen del ciclo contiene elementos que requieren atencion, o cuando +la IA fallo (RF-13). Nunca envia nada fuera del equipo (Regla 5). + +El runner es inyectable para poder testear sin reproducir audio. +""" + +from __future__ import annotations + +import subprocess +import sys +from typing import Callable + +from . import config + +import infra_paths + + +def _script_aviso(): + """Ruta del script de aviso sonoro (repo de negocio).""" + return infra_paths.biblioteca_repo() / "scripts" / "aviso_sonoro.py" + + +def _runner_default(tipo: str) -> None: + """Reproduce el aviso sonoro lanzando el script en un subproceso.""" + subprocess.run([sys.executable, str(_script_aviso()), tipo], check=False) + + +def avisar(resumen: dict, cfg: dict | None = None, _runner: Callable[[str], None] | None = None) -> bool: + """Emite el aviso por voz si procede. Devuelve True si aviso. + + - Si `aviso.voz` esta desactivado en config -> no hace nada. + - Avisa si `fallo_ia` o si hay elementos en `requiere_atencion`. + """ + cfg = cfg or config.cargar() + if not cfg.get("aviso", {}).get("voz", True): + return False + + fallo_ia = bool(resumen.get("fallo_ia")) + atencion = bool(resumen.get("requiere_atencion")) + if not (fallo_ia or atencion): + return False + + tipo = "fin" if atencion else "error" + runner = _runner or _runner_default + try: + runner(tipo) + except Exception: + pass + return True diff --git a/scripts/monitor_estado/ciclo.py b/scripts/monitor_estado/ciclo.py new file mode 100644 index 0000000..5b66da9 --- /dev/null +++ b/scripts/monitor_estado/ciclo.py @@ -0,0 +1,133 @@ +# -*- coding: utf-8 -*- +"""ciclo.py — Orquestador de un ciclo del monitor (RF-1..RF-18, RF-23, RF-24). + +Ejecuta la recogida de todas las fuentes, el archivado por fuente, el resumen IA, +el aviso por voz y el mantenimiento (poda, marca de ultimo ciclo). Aisla los errores +por fuente: una fuente caida no aborta el ciclo (RF-6). + +Las dependencias externas (fuentes, IA, aviso) son inyectables para poder testear. +""" + +from __future__ import annotations + +from datetime import date, datetime, timedelta + +from . import archivo, aviso, config, resumen, vistos +from .fuentes import correo as correo_mod +from .fuentes import jira as jira_mod +from .fuentes import telegram as telegram_mod + + +def _desde_jira(ahora: datetime) -> datetime: + """Instante desde el que buscar incidencias nuevas: ultimo ciclo o, si no, 24 h.""" + ultimo = archivo.leer_ultimo_ciclo() + inicio = ultimo.get("inicio") + if inicio: + try: + return datetime.fromisoformat(inicio) + except Exception: + pass + return ahora - timedelta(hours=24) + + +def ejecutar(cfg: dict | None = None, hoy: date | None = None, deps: dict | None = None) -> dict: + """Ejecuta un ciclo completo. Devuelve el resumen estructurado del ciclo.""" + cfg = cfg or config.cargar() + ahora = datetime.now() + hoy = hoy or ahora.date() + deps = deps or {} + + telegram_recoger = deps.get("telegram_recoger", telegram_mod.recoger_todas) + correo_recoger = deps.get("correo_recoger", correo_mod.recoger_recientes) + jira_cliente = deps.get("jira_cliente") + ia_fn = deps.get("ia") + avisar_fn = deps.get("avisar", aviso.avisar) + + novedades: dict[str, list] = {} + errores: list[tuple[str, str]] = [] + + # --- Telegram --- + try: + resultado, errs = telegram_recoger(cfg, hoy) + for nombre, msgs in resultado.items(): + archivo.escribe_dia( + f"telegram/{nombre}", hoy, msgs, + id_fn=lambda m: m["id"], render_fn=telegram_mod.render_mensaje, + titulo=f"{nombre} — {hoy.isoformat()}", + ) + nuevos = vistos.filtrar_nuevos(f"telegram_{nombre}", msgs, id_fn=lambda m: m["id"]) + if nuevos: + novedades[f"telegram/{nombre}"] = nuevos + errores += [(f"telegram/{n}", e) for n, e in errs] + except Exception as e: + errores.append(("telegram", str(e))) + + # --- Correo --- + try: + correos = correo_recoger(horas=cfg.get("correo", {}).get("horas", 24)) + archivo.escribe_dia("correo", hoy, correos, id_fn=lambda m: m["id"], render_fn=correo_mod.render_correo) + nuevos = vistos.filtrar_nuevos("correo", correos, id_fn=lambda m: m["id"]) + if nuevos: + novedades["correo"] = nuevos + except Exception as e: + errores.append(("correo", str(e))) + + # --- Jira (solo lectura) --- + try: + cliente = jira_cliente or jira_mod.construir_cliente(cfg) + ds_nuevos = jira_mod.ds_nuevos(cliente, _desde_jira(ahora)) + ds_sin_asignar = jira_mod.ds_sin_asignar(cliente) + todos = ds_nuevos + ds_sin_asignar + archivo.escribe_dia("jira", hoy, todos, id_fn=lambda m: m["key"], render_fn=jira_mod.render_jira) + nuevos = vistos.filtrar_nuevos("jira", todos, id_fn=lambda m: m["key"]) + if nuevos: + novedades["jira"] = nuevos + except Exception as e: + errores.append(("jira", str(e))) + + # --- Resumen IA + clasificacion --- + res = resumen.generar(novedades, errores, cfg=cfg, _ia=ia_fn) + archivo.escribe_resumen(res) + + # --- Aviso por voz --- + emitido = avisar_fn(res, cfg=cfg, _runner=deps.get("aviso_runner")) + res["aviso_emitido"] = emitido + + # --- Mantenimiento --- + archivo.poda(cfg.get("retencion_dias", 30), hoy) + archivo.escribe_ultimo_ciclo({ + "inicio": ahora.isoformat(), + "fin": datetime.now().isoformat(), + "ok": not errores, + "fuentes_con_novedad": res.get("fuentes_con_novedad", []), + "errores": res.get("errores", []), + }) + return res + + +def texto_estado() -> str: + """Texto legible del ultimo resumen y estado, para el comando `estado` (RF-18).""" + r = archivo.leer_ultimo_ciclo() + resumen_json = {} + ruta = archivo.MONITOR_DIR / "resumen_actual.json" + if ruta.exists(): + import json + try: + resumen_json = json.loads(ruta.read_text(encoding="utf-8")) + except Exception: + resumen_json = {} + lineas = [ + "── MONITOR DE ESTADO ──", + f"Datos en : {archivo.MONITOR_DIR}", + f"Ultimo ciclo : {r.get('inicio', 'nunca')} (ok={r.get('ok', '-')})", + ] + if r.get("errores"): + lineas.append(f"Errores : {len(r['errores'])}") + lineas.append("── Ultimo resumen ──") + lineas.append(resumen_json.get("resumen_md", "(sin resumen)") or "(sin resumen)") + atencion = resumen_json.get("requiere_atencion", []) + if atencion: + lineas.append("Requiere atencion:") + for a in atencion: + lineas.append(f" - [{a.get('fuente', '?')}] {a.get('texto', '')} ({a.get('motivo', '')})") + return "\n".join(lineas) diff --git a/scripts/monitor_estado/config.example.json b/scripts/monitor_estado/config.example.json new file mode 100644 index 0000000..201a381 --- /dev/null +++ b/scripts/monitor_estado/config.example.json @@ -0,0 +1,19 @@ +{ + "intervalo_minutos": 30, + "retencion_dias": 30, + "telegram": { + "sesion": "~/.telegram-agent/sessions/monitor", + "telefono": "+34657866417", + "fuentes": [ + { "nombre": "oscar_rodriguez", "id": 129020043 }, + { "nombre": "isa", "id": 16182959 }, + { "nombre": "daniel_garcia", "id": 8249301783 }, + { "nombre": "antonio_romero", "id": 880864397 }, + { "nombre": "desarrollo_software", "id": -1002693936231 } + ] + }, + "correo": { "cuenta": "jminguez@prolongo.es", "horas": 24 }, + "jira": { "proyecto": "DS" }, + "ia": { "base_url": "http://localhost:20128/v1", "modelo": "antigravity/gemini-2.5-flash-lite" }, + "aviso": { "voz": true } +} diff --git a/scripts/monitor_estado/config.py b/scripts/monitor_estado/config.py new file mode 100644 index 0000000..2448215 --- /dev/null +++ b/scripts/monitor_estado/config.py @@ -0,0 +1,117 @@ +# -*- coding: utf-8 -*- +"""config.py — Carga y valida la configuracion del monitor de estado. + +La configuracion vive en `INFRAESTRUCTURA/estado/monitor/config.json` (runtime, no +versionada). Si no existe, se crea a partir de la plantilla `config.example.json` +que va junto a este modulo, de modo que el monitor arranca siempre con valores validos. + +Las 5 fuentes de Telegram y sus identificadores se resolvieron desde la sesion real +de Telegram (ver plan 005). +""" + +from __future__ import annotations + +import copy +import json +import sys +from pathlib import Path + +# El modulo vive en INFRAESTRUCTURA/scripts/monitor_estado/; su padre es scripts/, +# donde esta infra_paths.py. +_SCRIPTS_DIR = Path(__file__).resolve().parent.parent +if str(_SCRIPTS_DIR) not in sys.path: + sys.path.insert(0, str(_SCRIPTS_DIR)) + +import infra_paths # noqa: E402 + +# Directorio raiz de datos runtime del monitor (gitignored). +MONITOR_DIR = infra_paths.ESTADO_DIR / "monitor" +CONFIG_PATH = MONITOR_DIR / "config.json" +EXAMPLE_PATH = Path(__file__).resolve().parent / "config.example.json" + + +def default_config() -> dict: + """Configuracion por defecto (equivale a la plantilla de ejemplo).""" + return { + "intervalo_minutos": 30, + "retencion_dias": 30, + "telegram": { + # Sesion Telethon DEDICADA (no reutiliza la del watcher, que esta bloqueada). + "sesion": "~/.telegram-agent/sessions/monitor", + "telefono": "+34657866417", + "fuentes": [ + {"nombre": "oscar_rodriguez", "id": 129020043}, + {"nombre": "isa", "id": 16182959}, + {"nombre": "daniel_garcia", "id": 8249301783}, + {"nombre": "antonio_romero", "id": 880864397}, + {"nombre": "desarrollo_software", "id": -1002693936231}, + ], + }, + "correo": {"cuenta": "jminguez@prolongo.es", "horas": 24}, + "jira": {"proyecto": "DS"}, + "ia": {"base_url": "http://localhost:20128/v1", "modelo": "antigravity/gemini-2.5-flash-lite"}, + "aviso": {"voz": True}, + } + + +def _merge(base: dict, override: dict) -> dict: + """Mezcla recursiva de diccionarios (override gana; no borra claves base).""" + result = copy.deepcopy(base) + for clave, valor in override.items(): + if isinstance(valor, dict) and isinstance(result.get(clave), dict): + result[clave] = _merge(result[clave], valor) + else: + result[clave] = valor + return result + + +def _escribir_plantilla(destino: Path) -> None: + """Escribe la config de ejemplo en `destino` (creando el directorio).""" + destino.parent.mkdir(parents=True, exist_ok=True) + with open(destino, "w", encoding="utf-8") as f: + json.dump(default_config(), f, indent=2, ensure_ascii=False) + + +def cargar(crear_si_falta: bool = True) -> dict: + """Carga la configuracion, creando `config.json` desde la plantilla si no existe. + + Devuelve siempre un dict completo (las claves que falten se toman por defecto). + """ + cfg = default_config() + if not CONFIG_PATH.exists(): + if not crear_si_falta: + return cfg + _escribir_plantilla(CONFIG_PATH) + try: + with open(CONFIG_PATH, "r", encoding="utf-8") as f: + usuario = json.load(f) + if isinstance(usuario, dict): + cfg = _merge(cfg, usuario) + except Exception as e: # config corrupta -> seguimos con los valores por defecto + print(f"[monitor_estado] Aviso: no se pudo leer {CONFIG_PATH}: {e}", file=sys.stderr) + return cfg + + +def fuentes_telegram(cfg: dict | None = None) -> list[dict]: + """Devuelve la lista de fuentes de Telegram configuradas.""" + cfg = cfg or cargar() + return list(cfg.get("telegram", {}).get("fuentes", [])) + + +def resumen_config(cfg: dict | None = None) -> str: + """Texto legible con las fuentes configuradas (para el CLI/estado).""" + cfg = cfg or cargar() + lineas = [ + f"Intervalo : {cfg.get('intervalo_minutos')} min", + f"Retencion : {cfg.get('retencion_dias')} dias", + f"Correo : {cfg.get('correo', {}).get('cuenta')} (ultimas {cfg.get('correo', {}).get('horas')} h)", + f"Jira : proyecto {cfg.get('jira', {}).get('proyecto')}", + "Telegram:", + ] + for f in fuentes_telegram(cfg): + lineas.append(f" - {f.get('nombre')} (id {f.get('id')})") + return "\n".join(lineas) + + +if __name__ == "__main__": + print(resumen_config()) diff --git a/scripts/monitor_estado/fuentes/__init__.py b/scripts/monitor_estado/fuentes/__init__.py new file mode 100644 index 0000000..f67cd96 --- /dev/null +++ b/scripts/monitor_estado/fuentes/__init__.py @@ -0,0 +1,2 @@ +# -*- coding: utf-8 -*- +"""Fuentes del monitor de estado (Telegram, correo, Jira).""" diff --git a/scripts/monitor_estado/fuentes/correo.py b/scripts/monitor_estado/fuentes/correo.py new file mode 100644 index 0000000..5ba53fc --- /dev/null +++ b/scripts/monitor_estado/fuentes/correo.py @@ -0,0 +1,73 @@ +# -*- coding: utf-8 -*- +"""correo.py — Fuente de correo de empresa del monitor (RF-2, RF-6, RF-25). + +Recoge los correos de las ultimas 24 h (o los de hoy) de `jminguez@prolongo.es` +usando la API de Gmail con el OAuth ya existente (repositorio de negocio). + +SOLO LECTURA: usa `messages().list()` y `messages().get()`; nunca marca correos +como leidos ni modifica nada (RF-25). +""" + +from __future__ import annotations + +import sys + +from .. import config # noqa: F401 (asegura sys.path con scripts/) + +import infra_paths # ya disponible en sys.path (lo anade config.py) + + +def _bootstrap_google() -> None: + """Deja importable el paquete `common` del repo de negocio (google/).""" + repo = infra_paths.bootstrap_biblioteca_paths() + google_dir = repo / "google" + if google_dir.exists() and str(google_dir) not in sys.path: + sys.path.insert(0, str(google_dir)) + + +def construir_servicio(): + """Devuelve el servicio de Gmail autenticado (OAuth existente).""" + _bootstrap_google() + import common # google/common.py + + return common.gmail_service() + + +def _consulta(horas: int) -> str: + """Consulta Gmail para la ventana pedida (granularidad en dias).""" + dias = max(1, int(horas) // 24) + return f"in:inbox newer_than:{dias}d" + + +def recoger_recientes(gmail=None, horas: int = 24, limite: int = 50) -> list[dict]: + """Devuelve los correos de la ventana, con {id, fecha, de, asunto, snippet}. + + No modifica el estado de los correos (no marca leido). + """ + gmail = gmail or construir_servicio() + res = gmail.users().messages().list(userId="me", q=_consulta(horas), maxResults=limite).execute() + salida: list[dict] = [] + for m in res.get("messages", []): + full = gmail.users().messages().get( + userId="me", id=m["id"], format="metadata", + metadataHeaders=["Subject", "From", "Date"], + ).execute() + headers = {h["name"].lower(): h["value"] for h in full.get("payload", {}).get("headers", [])} + salida.append({ + "id": m["id"], + "fecha": headers.get("date", ""), + "de": headers.get("from", "?"), + "asunto": headers.get("subject", "(sin asunto)"), + "snippet": full.get("snippet", ""), + }) + return salida + + +def render_correo(m: dict) -> str: + """Render Markdown de un correo para el fichero del dia.""" + asunto = m.get("asunto", "(sin asunto)") + de = m.get("de", "?") + fecha = m.get("fecha", "") + snippet = (m.get("snippet") or "").strip() + linea = f"- **{fecha}** — **{de}:** {asunto}" + return f"{linea}\n {snippet}" if snippet else linea diff --git a/scripts/monitor_estado/fuentes/jira.py b/scripts/monitor_estado/fuentes/jira.py new file mode 100644 index 0000000..a539d0f --- /dev/null +++ b/scripts/monitor_estado/fuentes/jira.py @@ -0,0 +1,94 @@ +# -*- coding: utf-8 -*- +"""jira.py — Fuente Jira del monitor (RF-3, RF-4, RF-16). + +Recoge del proyecto DS (Service Desk): + - incidencias NUEVAS creadas desde el ciclo anterior, + - tickets SIN ASIGNAR no cerrados ni resueltos. + +SOLO LECTURA (Regla 15): reutiliza el `JiraClient` de gestion_jira, que únicamente +expone consultas GET/POST de busqueda; nunca crea, modifica ni comenta. +""" + +from __future__ import annotations + +import sys +from datetime import datetime + +from .. import config # noqa: F401 (asegura sys.path con scripts/) + +import infra_paths + +# Estados que NO cuentan como "sin asignar pendiente de atencion". +_ESTADOS_EXCLUIDOS = ( + "'En Curso', 'In Progress', 'En desarrollo', 'Pasado a Programación Área', " + "'Cerrado', 'Resuelto', 'Cancelado', 'Closed', 'Resolved', 'Canceled'" +) + + +def _gestion_jira_dir(cfg: dict | None = None): + """Directorio del repo gestion_jira (donde vive jira_api.py).""" + cfg = cfg or config.cargar() + override = cfg.get("jira", {}).get("gestion_jira_dir") + if override: + return override + return str(infra_paths.INFRA_ROOT.parent / "gestion_jira") + + +def construir_cliente(cfg: dict | None = None): + """Devuelve un JiraClient (solo lectura).""" + directorio = _gestion_jira_dir(cfg) + if directorio not in sys.path: + sys.path.insert(0, directorio) + from jira_api import JiraClient + + return JiraClient() + + +def _iso_para_jql(desde) -> str: + """Normaliza una fecha a 'YYYY-MM-DD HH:MM' para el JQL.""" + if isinstance(desde, datetime): + return desde.strftime("%Y-%m-%d %H:%M") + return str(desde) + + +def _normalizar(issue: dict) -> dict: + """Convierte una issue de Jira en un registro legible.""" + key = issue.get("key", "") + f = issue.get("fields", {}) or {} + return { + "id": key, + "key": key, + "resumen": f.get("summary", ""), + "estado": (f.get("status") or {}).get("name", ""), + "creado": (f.get("created", "") or "")[:10], + "prioridad": (f.get("priority") or {}).get("name", ""), + "url": f"https://prolongo.atlassian.net/browse/{key}", + } + + +def ds_nuevos(cliente, desde) -> list[dict]: + """Incidencias DS creadas desde `desde` (inclusive).""" + jql = f'project = DS AND created >= "{_iso_para_jql(desde)}" ORDER BY created DESC' + res = cliente.search(jql, max_results=100) + return [_normalizar(i) for i in res.get("issues", [])] + + +def ds_sin_asignar(cliente) -> list[dict]: + """Tickets DS sin asignar, no cerrados ni resueltos.""" + jql = ( + "project = DS AND assignee = EMPTY " + f"AND status not in ({_ESTADOS_EXCLUIDOS}) " + "AND statusCategory != Done " + "ORDER BY created DESC" + ) + res = cliente.search(jql, max_results=100) + return [_normalizar(i) for i in res.get("issues", [])] + + +def render_jira(m: dict) -> str: + """Render Markdown de un ticket para el fichero del dia.""" + key = m.get("key", "?") + url = m.get("url", "") + estado = m.get("estado", "") + resumen = m.get("resumen", "") + return f"- **[{key}]({url})** ({estado}) — {resumen}" diff --git a/scripts/monitor_estado/fuentes/telegram.py b/scripts/monitor_estado/fuentes/telegram.py new file mode 100644 index 0000000..ff6f5af --- /dev/null +++ b/scripts/monitor_estado/fuentes/telegram.py @@ -0,0 +1,263 @@ +# -*- 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) diff --git a/scripts/monitor_estado/lock.py b/scripts/monitor_estado/lock.py new file mode 100644 index 0000000..dd17ac6 --- /dev/null +++ b/scripts/monitor_estado/lock.py @@ -0,0 +1,91 @@ +# -*- coding: utf-8 -*- +"""lock.py — Instancia unica del monitor (RF-22). + +Usa un fichero `.lock` en el directorio de datos del monitor con el PID y la hora de +arranque de la instancia viva. Antes de arrancar, si hay un PID vivo, se aborta; si el +lock es de un proceso muerto (apagado/crash), se recupera. + +Motivo: la tarea programada y el plugin de OpenCode pueden coincidir; este lock evita +que convivan dos monitores recogiendo a la vez. +""" + +from __future__ import annotations + +import json +import os +from datetime import datetime +from pathlib import Path + +from . import config + +LOCK_PATH: Path = config.MONITOR_DIR / ".lock" + +try: + import psutil # type: ignore +except Exception: # pragma: no cover - psutil esta disponible en el entorno + psutil = None + + +def pid_vivo(pid: int) -> bool: + """Devuelve True si el PID corresponde a un proceso vivo (no zombie).""" + if pid <= 0: + return False + if psutil is not None: + try: + return psutil.pid_exists(pid) and psutil.Process(pid).status() != psutil.STATUS_ZOMBIE + except Exception: + return False + # Fallback sin psutil. + try: + os.kill(pid, 0) + return True + except OSError: + return False + + +def _leer_pid() -> int | None: + """Lee el PID del lock actual, o None si no se puede.""" + try: + with open(LOCK_PATH, "r", encoding="utf-8") as f: + return int(json.load(f).get("pid", 0)) or None + except Exception: + return None + + +def _escribir_lock() -> None: + """Escribe el lock con el PID y la hora de arranque actuales.""" + LOCK_PATH.parent.mkdir(parents=True, exist_ok=True) + with open(LOCK_PATH, "w", encoding="utf-8") as f: + json.dump({"pid": os.getpid(), "arranque": datetime.now().isoformat()}, f, indent=2) + + +def adquirir() -> bool: + """Intenta adquirir el lock. Devuelve False si ya hay una instancia viva.""" + propio = os.getpid() + if LOCK_PATH.exists(): + pid = _leer_pid() + if pid and pid_vivo(pid): + # Viva (propia o ajena) -> no arrancar otra instancia. + return False + # Lock caducado (proceso muerto): se elimina y se reintenta. + try: + LOCK_PATH.unlink() + except OSError: + pass + _escribir_lock() + return True + + +def liberar() -> None: + """Elimina el lock solo si pertenece a nuestro propio PID.""" + try: + if LOCK_PATH.exists() and _leer_pid() == os.getpid(): + LOCK_PATH.unlink() + except OSError: + pass + + +def esta_bloqueado() -> bool: + """True si hay una instancia viva (distinta o igual) ocupando el lock.""" + pid = _leer_pid() if LOCK_PATH.exists() else None + return bool(pid and pid_vivo(pid)) diff --git a/scripts/monitor_estado/resumen.py b/scripts/monitor_estado/resumen.py new file mode 100644 index 0000000..a955685 --- /dev/null +++ b/scripts/monitor_estado/resumen.py @@ -0,0 +1,134 @@ +# -*- coding: utf-8 -*- +"""resumen.py — Resumen IA del estado y clasificacion "requiere atencion" (RF-11/12/13). + +Construye un prompt con las novedades del ciclo y llama al gateway IA ya configurado +en el entorno (OmniRoute, endpoint OpenAI-compatible en localhost). La IA devuelve un +JSON con el resumen en espanol y la lista de elementos que requieren atencion. + +Si la IA no esta disponible, no se pierde nada: se archiva igualmente lo recogido y se +deja constancia del fallo (`fallo_ia=True`) para que el orquestador avise (RF-13). + +El nucleo es testeable inyectando `_ia` (funcion prompt->texto) sin tocar la red. +""" + +from __future__ import annotations + +import json +import re +from datetime import datetime +from typing import Callable + +import requests + +from . import config + +_SYSTEM = ( + "Eres el asistente de estado de un jefe de desarrollo. Recibes las novedades del dia " + "(mensajes de Telegram, correos y tickets de Jira DS). Devuelve SOLO un JSON valido, sin " + "texto adicional, con esta forma: " + '{"resumen_md": "", ' + '"requiere_atencion": [{"fuente": "...", "quien": "...", "texto": "...", "motivo": "..."}]}. ' + "Marca en requiere_atencion solo incidencias, bloqueos, errores, paros, solicitudes urgentes, " + "tickets sin asignar o cambios de estado relevantes. Si no hay nada relevante, deja la lista vacia." +) + + +def _texto_novedad(fuente: str, item) -> str: + """Convierte un elemento de una fuente en una linea legible para el prompt.""" + if not isinstance(item, dict): + return f"[{fuente}] {item}" + if fuente.startswith("telegram"): + return f"[{fuente}] {item.get('quien', '?')}: {item.get('texto', '')}" + if fuente == "correo": + return f"[correo] {item.get('de', '?')}: {item.get('asunto', '')} — {item.get('snippet', '')}" + if fuente == "jira": + return f"[jira] {item.get('key', '?')} ({item.get('estado', '')}): {item.get('resumen', '')}" + return f"[{fuente}] {item}" + + +def construir_prompt(novedades: dict, errores: list) -> str: + """Texto del prompt a partir de las novedades y errores del ciclo.""" + lineas: list[str] = [] + for fuente, items in novedades.items(): + if not items: + continue + lineas.append(f"## {fuente}") + for it in items: + lineas.append("- " + _texto_novedad(fuente, it)) + if errores: + lineas.append("## errores de recogida") + for fuente, err in errores: + lineas.append(f"- [{fuente}] {err}") + return "\n".join(lineas) + + +def _extraer_json(texto: str) -> dict: + """Extrae el primer objeto JSON del texto (tolerando fences de markdown).""" + limpio = re.sub(r"^```(?:json)?|```$", "", texto.strip(), flags=re.MULTILINE).strip() + inicio = limpio.find("{") + fin = limpio.rfind("}") + if inicio == -1 or fin == -1: + raise ValueError("la IA no devolvio JSON") + return json.loads(limpio[inicio:fin + 1]) + + +def _llamar_ia(prompt: str, cfg: dict) -> str: + """Llama al gateway IA (OmniRoute, OpenAI-compatible) y devuelve el texto.""" + ia = cfg.get("ia", {}) + base_url = ia.get("base_url", "http://localhost:20128/v1").rstrip("/") + modelo = ia.get("modelo", "auto/fast") + headers = {"Content-Type": "application/json"} + api_key = ia.get("api_key") + if api_key: + headers["Authorization"] = f"Bearer {api_key}" + body = { + "model": modelo, + "messages": [ + {"role": "system", "content": _SYSTEM}, + {"role": "user", "content": prompt}, + ], + "temperature": 0.2, + "stream": False, + } + r = requests.post(f"{base_url}/chat/completions", json=body, headers=headers, timeout=90) + r.raise_for_status() + return r.json()["choices"][0]["message"]["content"] + + +def generar( + novedades: dict, + errores: list, + cfg: dict | None = None, + _ia: Callable[[str], str] | None = None, +) -> dict: + """Genera el resumen del ciclo. + + Devuelve un dict con: generado, fuentes_con_novedad, errores, fallo_ia, + requiere_atencion[] y resumen_md. + """ + cfg = cfg or config.cargar() + fuentes_con_novedad = [k for k, v in novedades.items() if v] + resultado = { + "generado": datetime.now().isoformat(), + "fuentes_con_novedad": fuentes_con_novedad, + "errores": [{"fuente": f, "error": e} for f, e in errores], + "fallo_ia": False, + "requiere_atencion": [], + "resumen_md": "", + } + + if not fuentes_con_novedad and not errores: + resultado["resumen_md"] = "Sin novedades." + return resultado + + prompt = construir_prompt(novedades, errores) + try: + texto = _ia(prompt) if _ia is not None else _llamar_ia(prompt, cfg) + datos = _extraer_json(texto) + resultado["resumen_md"] = datos.get("resumen_md", "") or "" + atencion = datos.get("requiere_atencion", []) or [] + resultado["requiere_atencion"] = atencion if isinstance(atencion, list) else [] + except Exception as e: # RF-13: no perder lo recogido y dejar constancia + resultado["fallo_ia"] = True + resultado["resumen_md"] = f"No se pudo generar el resumen IA: {e}" + return resultado diff --git a/scripts/monitor_estado/vistos.py b/scripts/monitor_estado/vistos.py new file mode 100644 index 0000000..25b9b51 --- /dev/null +++ b/scripts/monitor_estado/vistos.py @@ -0,0 +1,94 @@ +# -*- coding: utf-8 -*- +"""vistos.py — Estado persistente de "ya visto" por fuente (RF-8, RF-23). + +Guarda, por cada fuente, el conjunto de identificadores ya procesados para no +volver a contarlos ni a avisar por ellos en el siguiente ciclo. El estado se +persiste en `estado/monitor/visto/.json` y sobrevive a reinicios. + +Los identificadores pueden ser enteros (Telegram: id de mensaje) o cadenas +(correo: id de mensaje de Gmail; Jira: clave del ticket). Internamente se +normalizan a texto. +""" + +from __future__ import annotations + +import json +import re +from pathlib import Path +from typing import Callable, Iterable + +from . import config + +VISTO_DIR: Path = config.MONITOR_DIR / "visto" + +# Limite de ids guardados por fuente (ventana movil; evita ficheros infinitos). +MAX_IDS = 2000 + + +def _nombre_fichero(clave: str) -> str: + """Nombre de fichero seguro para la clave de la fuente.""" + seguro = re.sub(r"[^A-Za-z0-9_.-]", "_", clave) + return f"{seguro}.json" + + +def _ruta(clave: str) -> Path: + return VISTO_DIR / _nombre_fichero(clave) + + +def _cargar(clave: str) -> dict: + """Carga el estado de la fuente (lista de ids vistos).""" + ruta = _ruta(clave) + if ruta.exists(): + try: + with open(ruta, "r", encoding="utf-8") as f: + data = json.load(f) + if isinstance(data, dict): + return data + except Exception: + pass + return {"ids_vistos": []} + + +def _guardar(clave: str, data: dict) -> None: + """Persiste el estado de la fuente (recortando a MAX_IDS).""" + VISTO_DIR.mkdir(parents=True, exist_ok=True) + ids = data.get("ids_vistos", []) + data["ids_vistos"] = ids[-MAX_IDS:] + with open(_ruta(clave), "w", encoding="utf-8") as f: + json.dump(data, f, indent=2, ensure_ascii=False) + + +def _norm(x) -> str: + return str(x) + + +def marcar(clave: str, ids: Iterable) -> None: + """Añade ids al estado de la fuente y persiste.""" + data = _cargar(clave) + vistos = set(data.get("ids_vistos", [])) + nuevos = {_norm(i) for i in ids} - vistos + if not nuevos: + return + data["ids_vistos"] = list(vistos | nuevos) + _guardar(clave, data) + + +def filtrar_nuevos( + clave: str, + elementos: Iterable, + id_fn: Callable = lambda e: e, + marcar_automatico: bool = True, +) -> list: + """Devuelve los elementos cuyo id no estaba visto. + + Si `marcar_automatico` es True (por defecto), los ids nuevos quedan marcados + en el mismo paso, de modo que el siguiente ciclo no los repita. + """ + elementos = list(elementos) + vistos = set(_cargar(clave).get("ids_vistos", [])) + nuevos = [e for e in elementos if _norm(id_fn(e)) not in vistos] + if marcar_automatico and nuevos: + data = _cargar(clave) + data["ids_vistos"] = list(vistos | {_norm(id_fn(e)) for e in nuevos}) + _guardar(clave, data) + return nuevos diff --git a/scripts/monitor_estado_watchdog.ps1 b/scripts/monitor_estado_watchdog.ps1 new file mode 100644 index 0000000..1b86cfe --- /dev/null +++ b/scripts/monitor_estado_watchdog.ps1 @@ -0,0 +1,59 @@ +# monitor_estado_watchdog.ps1 — Asegura que el Monitor de estado esta corriendo (una instancia). +# +# Uso: +# powershell -ExecutionPolicy Bypass -File monitor_estado_watchdog.ps1 -Once # lanza si no corre +# powershell -ExecutionPolicy Bypass -File monitor_estado_watchdog.ps1 -Stop # detiene el monitor +# +# Lo llaman: la tarea programada al iniciar sesion y el plugin global de OpenCode +# (monitor-estado-bootstrap.ts). El propio monitor tiene ademas un lock (lock.py), +# asi que aunque ambos coincidan no habra dos instancias. + +param( + [switch]$Once, + [switch]$Stop +) + +$ErrorActionPreference = 'SilentlyContinue' + +$ScriptsDir = $PSScriptRoot +$Repo = Split-Path -Parent $ScriptsDir +$Script = Join-Path $ScriptsDir 'monitor_estado.py' +$LogDir = Join-Path $Repo 'estado\monitor' +$OutLog = Join-Path $LogDir 'monitor_estado.out.log' +$ErrLog = Join-Path $LogDir 'monitor_estado.err.log' +$LockFile = Join-Path $LogDir '.lock' + +if (-not (Test-Path $LogDir)) { New-Item -ItemType Directory -Path $LogDir -Force | Out-Null } + +function Get-MonitorProcesos { + Get-CimInstance Win32_Process | + Where-Object { $_.CommandLine -match 'monitor_estado\.py\s+daemon' } +} + +if ($Stop) { + $procs = Get-MonitorProcesos + foreach ($p in $procs) { Stop-Process -Id $p.ProcessId -Force } + Remove-Item -LiteralPath $LockFile -Force -ErrorAction SilentlyContinue + Write-Output "[monitor_estado] Detenido ($($procs.Count) proceso(s))." + exit 0 +} + +$procs = Get-MonitorProcesos +if ($procs) { + Write-Output "[monitor_estado] Ya en marcha ($($procs.Count) proceso). Nada que hacer." + exit 0 +} + +# Lanzar el daemon en segundo plano, oculto y con log. +$python = (Get-Command python).Source +if (-not $python) { Write-Output "[monitor_estado] Python no encontrado en PATH."; exit 1 } + +Start-Process -FilePath $python ` + -ArgumentList @('-u', $Script, 'daemon') ` + -WorkingDirectory $Repo ` + -WindowStyle Hidden ` + -RedirectStandardOutput $OutLog ` + -RedirectStandardError $ErrLog + +Write-Output "[monitor_estado] Daemon lanzado (log: $OutLog)." +exit 0 diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..67ecd64 --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,13 @@ +# -*- coding: utf-8 -*- +"""conftest.py — Deja importable el paquete `monitor_estado` en los tests. + +Anade INFRAESTRUCTURA/scripts al sys.path para poder hacer +`import monitor_estado.` sin instalar el paquete. +""" + +import sys +from pathlib import Path + +SCRIPTS_DIR = Path(__file__).resolve().parent.parent / "scripts" +if str(SCRIPTS_DIR) not in sys.path: + sys.path.insert(0, str(SCRIPTS_DIR)) diff --git a/tests/test_archivo.py b/tests/test_archivo.py new file mode 100644 index 0000000..cf05655 --- /dev/null +++ b/tests/test_archivo.py @@ -0,0 +1,73 @@ +# -*- coding: utf-8 -*- +"""Tests de `monitor_estado.archivo` — ficheros del dia, reescritura y poda (RF-5,7,9,10,18).""" + +from datetime import date, timedelta + +import monitor_estado.archivo as archivo + + +def _render(m): + return f"- **{m['id']}**: {m['t']}" + + +def test_archivo_reescribe_sin_duplicados(tmp_path, monkeypatch): + """El fichero del dia contiene todo lo de hoy y nada duplicado (RF-9).""" + monkeypatch.setattr(archivo, "MONITOR_DIR", tmp_path) + hoy = date(2026, 10, 3) + + p = archivo.escribe_dia( + "telegram/test", hoy, + [{"id": 1, "t": "a"}, {"id": 2, "t": "b"}, {"id": 2, "t": "b"}], + id_fn=lambda m: m["id"], render_fn=_render, + ) + contenido = p.read_text(encoding="utf-8") + assert contenido.count("- **1**") == 1 + assert contenido.count("- **2**") == 1 + + # Nuevo ciclo con lo del dia completo: el fichero se reescribe con el estado actual. + archivo.escribe_dia( + "telegram/test", hoy, + [{"id": 1, "t": "a"}, {"id": 2, "t": "b"}, {"id": 3, "t": "c"}], + id_fn=lambda m: m["id"], render_fn=_render, + ) + contenido = p.read_text(encoding="utf-8") + assert "- **3**" in contenido + assert contenido.count("- **1**") == 1 + + +def test_archivo_poda_30_dias_y_conserva_recientes(tmp_path, monkeypatch): + """Se borran los ficheros de mas de 30 dias y se conservan los recientes (RF-10).""" + monkeypatch.setattr(archivo, "MONITOR_DIR", tmp_path) + hoy = date(2026, 10, 3) + reciente = tmp_path / "telegram" / "test" + reciente.mkdir(parents=True) + f_hoy = reciente / f"{hoy.isoformat()}.md" + f_10 = reciente / f"{(hoy - timedelta(days=10)).isoformat()}.md" + f_40 = reciente / f"{(hoy - timedelta(days=40)).isoformat()}.md" + for f in (f_hoy, f_10, f_40): + f.write_text("x", encoding="utf-8") + + borrados = archivo.poda(retencion_dias=30, hoy=hoy) + + assert f_hoy.exists() and f_10.exists() + assert not f_40.exists() + assert f_40 in borrados + + +def test_escribe_resumen(tmp_path, monkeypatch): + """El resumen se guarda en md y json consultables (RF-18).""" + monkeypatch.setattr(archivo, "MONITOR_DIR", tmp_path) + resumen = {"resumen_md": "Todo ok", "requiere_atencion": [], "errores": []} + archivo.escribe_resumen(resumen) + + assert (tmp_path / "resumen_actual.md").read_text(encoding="utf-8") == "Todo ok" + assert '"resumen_md": "Todo ok"' in (tmp_path / "resumen_actual.json").read_text(encoding="utf-8") + + +def test_escribe_dia_de_fuente_sin_novedad(tmp_path, monkeypatch): + """Una fuente sin elementos genera un fichero vacio (no falla) (RF-5).""" + monkeypatch.setattr(archivo, "MONITOR_DIR", tmp_path) + p = archivo.escribe_dia( + "correo", date(2026, 10, 3), [], id_fn=lambda m: m, render_fn=lambda m: str(m) + ) + assert p.exists() diff --git a/tests/test_aviso.py b/tests/test_aviso.py new file mode 100644 index 0000000..b2b62d4 --- /dev/null +++ b/tests/test_aviso.py @@ -0,0 +1,32 @@ +# -*- coding: utf-8 -*- +"""Tests de `monitor_estado.aviso` — voz solo si requiere atencion (RF-14, RF-15).""" + +import monitor_estado.aviso as aviso + + +def test_aviso_dispara_con_atencion(): + llamadas = [] + r = {"fallo_ia": False, "requiere_atencion": [{"fuente": "jira", "texto": "caido"}]} + assert aviso.avisar(r, cfg={"aviso": {"voz": True}}, _runner=lambda t: llamadas.append(t)) is True + assert llamadas == ["fin"] + + +def test_ciclo_sin_novedad_no_avisa(): + llamadas = [] + r = {"fallo_ia": False, "requiere_atencion": [], "resumen_md": "Sin novedades."} + assert aviso.avisar(r, cfg={"aviso": {"voz": True}}, _runner=lambda t: llamadas.append(t)) is False + assert llamadas == [] + + +def test_aviso_dispara_con_fallo_ia(): + llamadas = [] + r = {"fallo_ia": True, "requiere_atencion": []} + assert aviso.avisar(r, cfg={"aviso": {"voz": True}}, _runner=lambda t: llamadas.append(t)) is True + assert llamadas == ["error"] + + +def test_aviso_desactivado_por_config(): + llamadas = [] + r = {"fallo_ia": False, "requiere_atencion": [{"texto": "x"}]} + assert aviso.avisar(r, cfg={"aviso": {"voz": False}}, _runner=lambda t: llamadas.append(t)) is False + assert llamadas == [] diff --git a/tests/test_ciclo.py b/tests/test_ciclo.py new file mode 100644 index 0000000..4a79c02 --- /dev/null +++ b/tests/test_ciclo.py @@ -0,0 +1,81 @@ +# -*- coding: utf-8 -*- +"""Tests de `monitor_estado.ciclo` — aislamiento de fuentes y ciclo completo (RF-6, RF-9).""" + +from datetime import date + +import monitor_estado.archivo as archivo +import monitor_estado.ciclo as ciclo +import monitor_estado.vistos as vistos + + +class _EmptyJira: + def search(self, jql, fields=None, max_results=50): + return {"issues": []} + + +def _preparar_tmp(tmp_path, monkeypatch): + monkeypatch.setattr(archivo, "MONITOR_DIR", tmp_path) + monkeypatch.setattr(vistos, "VISTO_DIR", tmp_path / "visto") + + +def test_aislamiento_fuente_fallida_continua(tmp_path, monkeypatch): + """Si todas las fuentes fallan, el ciclo termina y registra los errores (RF-6).""" + _preparar_tmp(tmp_path, monkeypatch) + + def bad_tg(cfg, hoy): + raise RuntimeError("telegram caido") + + def bad_correo(horas=24, limite=50): + raise RuntimeError("gmail caido") + + class BadJira: + def search(self, *a, **k): + raise RuntimeError("jira caido") + + res = ciclo.ejecutar( + cfg={"retencion_dias": 30, "correo": {"horas": 24}}, + deps={ + "telegram_recoger": bad_tg, + "correo_recoger": bad_correo, + "jira_cliente": BadJira(), + "ia": lambda p: '{"resumen_md": "x", "requiere_atencion": []}', + "avisar": lambda r, cfg=None, _runner=None: False, + }, + ) + assert len(res["errores"]) == 3 + assert res["fallo_ia"] is False + assert (tmp_path / "ultimo_ciclo.json").exists() + + +def test_ciclo_archiva_y_avisa(tmp_path, monkeypatch): + """Un ciclo con novedad archiva el dia y avisa si hay atencion (RF-9, RF-14).""" + _preparar_tmp(tmp_path, monkeypatch) + emitidos = [] + + def tg(cfg, hoy): + return ({"oscar_rodriguez": [{"id": 1, "fecha": "2026-10-03T09:00:00", "quien": "Oscar", "texto": "incidencia"}]}, []) + + def correo(horas=24, limite=50): + return [] + + def ia(prompt): + return '{"resumen_md": "hay incidencia", "requiere_atencion": [{"fuente": "telegram", "texto": "incidencia"}]}' + + def avisar(res, cfg=None, _runner=None): + emitidos.append(res.get("requiere_atencion")) + return True + + res = ciclo.ejecutar( + cfg={"retencion_dias": 30, "correo": {"horas": 24}}, + hoy=date(2026, 10, 3), + deps={ + "telegram_recoger": tg, + "correo_recoger": correo, + "jira_cliente": _EmptyJira(), + "ia": ia, + "avisar": avisar, + }, + ) + assert (tmp_path / "telegram" / "oscar_rodriguez" / "2026-10-03.md").exists() + assert res["aviso_emitido"] is True + assert emitidos and emitidos[0] diff --git a/tests/test_correo.py b/tests/test_correo.py new file mode 100644 index 0000000..e0c76df --- /dev/null +++ b/tests/test_correo.py @@ -0,0 +1,74 @@ +# -*- coding: utf-8 -*- +"""Tests de `monitor_estado.fuentes.correo` — ventana 24h y solo lectura (RF-2, RF-25).""" + +import monitor_estado.fuentes.correo as correo + + +class _Req: + def __init__(self, resp): + self._resp = resp + + def execute(self): + return self._resp + + +class _Messages: + def __init__(self): + self.queries = [] + self.gets = [] + self.modificados = [] + + def list(self, userId, q, maxResults): + self.queries.append(q) + return _Req({"messages": [{"id": "m1"}, {"id": "m2"}]}) + + def get(self, userId, id, format, metadataHeaders): + self.gets.append(id) + return _Req({ + "snippet": f"resumen de {id}", + "payload": {"headers": [ + {"name": "Subject", "value": f"Asunto {id}"}, + {"name": "From", "value": "Remitente "}, + {"name": "Date", "value": "Sat, 3 Oct 2026 09:00:00 +0200"}, + ]}, + }) + + def modify(self, **kwargs): # no debe llamarse nunca (RF-25) + self.modificados.append(kwargs) + return _Req({}) + + +class _Users: + def __init__(self): + self._m = _Messages() + + def messages(self): + return self._m + + +class FakeGmail: + def __init__(self): + self._u = _Users() + + def users(self): + return self._u + + +def test_correo_ventana_24h(): + gmail = FakeGmail() + correos = correo.recoger_recientes(gmail, horas=24) + assert len(correos) == 2 + assert correos[0]["id"] == "m1" + assert correos[0]["asunto"] == "Asunto m1" + assert "newer_than:1d" in gmail.users().messages().queries[0] + + +def test_correo_no_modifica_estado(): + gmail = FakeGmail() + correo.recoger_recientes(gmail, horas=24) + assert gmail.users().messages().modificados == [] # no marca leido (RF-25) + + +def test_render_correo(): + txt = correo.render_correo({"fecha": "Sat, 3 Oct 2026 09:00:00 +0200", "de": "a@b.com", "asunto": "Hola", "snippet": "cuerpo"}) + assert "a@b.com" in txt and "Hola" in txt diff --git a/tests/test_jira.py b/tests/test_jira.py new file mode 100644 index 0000000..859b8b9 --- /dev/null +++ b/tests/test_jira.py @@ -0,0 +1,58 @@ +# -*- coding: utf-8 -*- +"""Tests de `monitor_estado.fuentes.jira` — DS nuevos y sin asignar, solo GET (RF-3, RF-4, RF-16).""" + +import monitor_estado.fuentes.jira as jira + + +def _issue(key, summary, status="Pendiente", created="2026-10-03T08:00:00.000+0200"): + return {"key": key, "fields": { + "summary": summary, + "status": {"name": status}, + "created": created, + "priority": {"name": "Media"}, + }} + + +class FakeJira: + """Solo implementa `search` (lectura). Si el codigo intentara escribir, fallaria.""" + + def __init__(self): + self.calls = [] + + def search(self, jql, fields=None, max_results=50): + self.calls.append(jql) + if "assignee = EMPTY" in jql: + issues = [_issue("DS-100", "Ticket sin asignar")] + elif "created >=" in jql: + issues = [_issue("DS-200", "Ticket nuevo")] + else: + issues = [] + return {"issues": issues} + + +def test_jira_jql_ds_sin_asignar(): + fake = FakeJira() + out = jira.ds_sin_asignar(fake) + assert any("assignee = EMPTY" in c and "project = DS" in c for c in fake.calls) + assert out[0]["key"] == "DS-100" + assert out[0]["estado"] == "Pendiente" + + +def test_jira_ds_nuevos_ventana(): + fake = FakeJira() + out = jira.ds_nuevos(fake, "2026-10-03 00:00") + assert "created >=" in fake.calls[-1] + assert out[0]["key"] == "DS-200" + + +def test_jira_solo_lectura(): + """El cliente usado solo expone `search` (no hay escritura) — RF-16.""" + fake = FakeJira() + jira.ds_sin_asignar(fake) + jira.ds_nuevos(fake, "2026-10-03 00:00") + assert not hasattr(fake, "create_issue") and not hasattr(fake, "post") + + +def test_render_jira(): + txt = jira.render_jira({"key": "DS-1", "resumen": "Algo", "estado": "Pendiente", "url": "u"}) + assert "DS-1" in txt and "Algo" in txt and "Pendiente" in txt diff --git a/tests/test_lock.py b/tests/test_lock.py new file mode 100644 index 0000000..5342bf2 --- /dev/null +++ b/tests/test_lock.py @@ -0,0 +1,40 @@ +# -*- coding: utf-8 -*- +"""Tests de `monitor_estado.lock` — instancia unica (RF-22).""" + +import json + +import monitor_estado.lock as lock + + +def test_lock_impide_segunda_instancia(tmp_path, monkeypatch): + """Una segunda adquisicion con la primera viva debe fallar (RF-22).""" + monkeypatch.setattr(lock, "LOCK_PATH", tmp_path / ".lock") + + assert lock.adquirir() is True + # Misma PID viva -> la segunda instancia no arranca. + assert lock.adquirir() is False + + lock.liberar() + assert not (tmp_path / ".lock").exists() + assert lock.adquirir() is True + lock.liberar() + + +def test_lock_recupera_lock_caducado(tmp_path, monkeypatch): + """Un lock de un PID muerto no debe bloquear (se recupera).""" + ruta = tmp_path / ".lock" + monkeypatch.setattr(lock, "LOCK_PATH", ruta) + ruta.write_text(json.dumps({"pid": 999999, "arranque": "2026-01-01T00:00:00"}), encoding="utf-8") + + assert lock.adquirir() is True + lock.liberar() + + +def test_lock_liberar_no_borra_lock_ajeno(tmp_path, monkeypatch): + """`liberar` solo elimina el lock si es de nuestro propio PID.""" + ruta = tmp_path / ".lock" + monkeypatch.setattr(lock, "LOCK_PATH", ruta) + ruta.write_text(json.dumps({"pid": 999999, "arranque": "x"}), encoding="utf-8") + + lock.liberar() + assert ruta.exists() # no era nuestro -> se respeta diff --git a/tests/test_resumen.py b/tests/test_resumen.py new file mode 100644 index 0000000..165e598 --- /dev/null +++ b/tests/test_resumen.py @@ -0,0 +1,45 @@ +# -*- coding: utf-8 -*- +"""Tests de `monitor_estado.resumen` — resumen IA y clasificacion (RF-11, RF-12, RF-13).""" + +import monitor_estado.resumen as resumen + + +def test_resumen_clasifica_atencion(): + def fake_ia(prompt): + return ('{"resumen_md": "Hay una incidencia de SQL.", ' + '"requiere_atencion": [{"fuente": "jira", "quien": "DS-3648", ' + '"texto": "PC526 sin conexion", "motivo": "ticket sin asignar"}]}') + + r = resumen.generar({"jira": [{"key": "DS-3648", "resumen": "PC526 sin conexion"}]}, [], _ia=fake_ia) + assert r["fallo_ia"] is False + assert len(r["requiere_atencion"]) == 1 + assert r["requiere_atencion"][0]["quien"] == "DS-3648" + assert "incidencia" in r["resumen_md"].lower() + + +def test_resumen_fallo_ia_deja_constancia(): + def bad_ia(prompt): + raise RuntimeError("OmniRoute caido") + + r = resumen.generar({"jira": [{"key": "DS-1", "resumen": "x"}]}, [("correo", "token caducado")], _ia=bad_ia) + assert r["fallo_ia"] is True + assert r["requiere_atencion"] == [] + assert r["resumen_md"] # deja constancia del fallo + assert r["errores"] # conserva los errores del ciclo + + +def test_resumen_parsea_json_con_fences(): + def fake_ia(prompt): + return '```json\n{"resumen_md": "ok", "requiere_atencion": []}\n```' + + r = resumen.generar({"correo": [{"asunto": "hola"}]}, [], _ia=fake_ia) + assert r["fallo_ia"] is False and r["resumen_md"] == "ok" + + +def test_resumen_sin_novedades_no_llama_ia(): + def no_llamar(prompt): + raise AssertionError("no deberia llamarse a la IA sin novedades") + + r = resumen.generar({}, [], _ia=no_llamar) + assert r["fallo_ia"] is False + assert r["requiere_atencion"] == [] diff --git a/tests/test_telegram.py b/tests/test_telegram.py new file mode 100644 index 0000000..66d7150 --- /dev/null +++ b/tests/test_telegram.py @@ -0,0 +1,91 @@ +# -*- coding: utf-8 -*- +"""Tests de `monitor_estado.fuentes.telegram` — nucleo y aislamiento (RF-1, RF-6).""" + +from datetime import date, datetime, timezone +from types import SimpleNamespace + +import monitor_estado.fuentes.telegram as tg + + +class FakeMsg: + """Mensaje crudo minimo compatible con lo que espera normalizar/recoger.""" + + def __init__(self, id, dt, text, sender_id=1, sender=None): + self.id = id + self.date = dt + self.message = text + self.sender_id = sender_id + self._sender_name = sender + self._sender = sender or "Anon" + + async def get_sender(self): + return SimpleNamespace(first_name=self._sender) + + +class FakeClient: + """Cliente Telethon falso (sin red).""" + + def __init__(self, por_id, malos=()): + self.por_id = por_id + self.malos = set(malos) + self.conectado = False + + async def connect(self): + self.conectado = True + + async def is_user_authorized(self): + return True + + async def disconnect(self): + self.conectado = False + + async def get_dialogs(self, limit=None): + return [] + + async def get_entity(self, cid): + if cid in self.malos: + raise RuntimeError("sin acceso a la fuente") + return cid + + def iter_messages(self, entidad, limit=None): + async def gen(): + for m in self.por_id.get(entidad, []): + yield m + return gen() + + +def test_normalizar_filtra_dia_y_ordena(): + dia = date(2026, 10, 3) + m1 = FakeMsg(2, datetime(2026, 10, 3, 15, 0, tzinfo=timezone.utc), "segundo", sender="Ana") + m2 = FakeMsg(1, datetime(2026, 10, 3, 9, 0, tzinfo=timezone.utc), "primero", sender="Ana") + m3 = FakeMsg(3, datetime(2026, 10, 2, 9, 0, tzinfo=timezone.utc), "ayer", sender="Ana") + out = tg.normalizar_mensajes([m1, m2, m3], dia, "test") + assert [r["id"] for r in out] == [1, 2] # solo hoy y ordenado + assert out[0]["quien"] == "Ana" and out[0]["fuente"] == "test" + + +def test_recoger_fuente_para_en_el_dia(): + dia = date(2026, 10, 3) + hoy = FakeMsg(5, datetime(2026, 10, 3, 12, 0, tzinfo=timezone.utc), "hoy", sender="Ana") + ayer = FakeMsg(4, datetime(2026, 10, 2, 12, 0, tzinfo=timezone.utc), "ayer", sender="Ana") + client = FakeClient({7: [hoy, ayer]}) + import asyncio + out = asyncio.run(tg.recoger_fuente(client, "oscar", 7, dia)) + assert [r["id"] for r in out] == [5] + + +def test_recoger_todas_aisla_fuente_caida(): + dia = date(2026, 10, 3) + ok = FakeMsg(1, datetime(2026, 10, 3, 12, 0, tzinfo=timezone.utc), "hola", sender="Ana") + cfg = {"telegram": {"fuentes": [{"nombre": "buena", "id": 10}, {"nombre": "mala", "id": 20}]}} + client = FakeClient({10: [ok]}, malos={20}) + + resultado, errores = tg.recoger_todas(cfg, dia, client_factory=lambda c: client) + assert "buena" in resultado and len(resultado["buena"]) == 1 + assert len(errores) == 1 and errores[0][0] == "mala" + + +def test_render_mensaje(): + m = {"fecha": "2026-10-03T09:12:00+02:00", "quien": "Oscar", "texto": "linea1\nlinea2"} + txt = tg.render_mensaje(m) + assert "[09:12] Oscar" in txt and "linea1" in txt and "linea2" in txt diff --git a/tests/test_vistos.py b/tests/test_vistos.py new file mode 100644 index 0000000..a7369c8 --- /dev/null +++ b/tests/test_vistos.py @@ -0,0 +1,48 @@ +# -*- coding: utf-8 -*- +"""Tests de `monitor_estado.vistos` — dedup persistente (RF-8, RF-23).""" + +import monitor_estado.vistos as vistos + + +def test_vistos_filtra_nuevos(tmp_path, monkeypatch): + """Lo ya visto no se vuelve a contar; lo nuevo si (RF-8).""" + monkeypatch.setattr(vistos, "VISTO_DIR", tmp_path) + clave = "telegram_test" + msgs = [{"id": 1, "t": "a"}, {"id": 2, "t": "b"}] + + nuevos = vistos.filtrar_nuevos(clave, msgs, id_fn=lambda m: m["id"]) + assert [m["id"] for m in nuevos] == [1, 2] + + # Segundo ciclo con los mismos mensajes -> nada nuevo. + assert vistos.filtrar_nuevos(clave, msgs, id_fn=lambda m: m["id"]) == [] + + # Llega uno nuevo -> solo ese. + mas = msgs + [{"id": 3, "t": "c"}] + assert [m["id"] for m in vistos.filtrar_nuevos(clave, mas, id_fn=lambda m: m["id"])] == [3] + + +def test_vistos_persiste_entre_llamadas(tmp_path, monkeypatch): + """El estado se guarda en disco y sobrevive a una nueva lectura (RF-23).""" + monkeypatch.setattr(vistos, "VISTO_DIR", tmp_path) + clave = "correo" + vistos.filtrar_nuevos(clave, ["m1", "m2"], id_fn=lambda x: x) + + assert (tmp_path / f"{clave}.json").exists() + # Nueva lectura del fichero: m1/m2 ya no son nuevos. + assert vistos.filtrar_nuevos(clave, ["m1", "m2"], id_fn=lambda x: x) == [] + + +def test_vistos_marcar_explicito(tmp_path, monkeypatch): + """`marcar` añade ids sin devolver nada.""" + monkeypatch.setattr(vistos, "VISTO_DIR", tmp_path) + vistos.marcar("jira", ["DS-1", "DS-2"]) + assert vistos.filtrar_nuevos("jira", ["DS-1"], id_fn=lambda x: x) == [] + assert vistos.filtrar_nuevos("jira", ["DS-3"], id_fn=lambda x: x) == ["DS-3"] + + +def test_vistos_idempotente_sin_marcar(tmp_path, monkeypatch): + """Con marcar_automatico=False, filtrar no altera el estado.""" + monkeypatch.setattr(vistos, "VISTO_DIR", tmp_path) + a = vistos.filtrar_nuevos("k", [1, 2], id_fn=lambda x: x, marcar_automatico=False) + b = vistos.filtrar_nuevos("k", [1, 2], id_fn=lambda x: x, marcar_automatico=False) + assert a == [1, 2] and b == [1, 2]