03-oct: superguardado automático
This commit is contained in:
58
scripts/monitor_estado/README.md
Normal file
58
scripts/monitor_estado/README.md
Normal file
@@ -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 (<YYYY-MM-DD>.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).
|
||||
11
scripts/monitor_estado/__init__.py
Normal file
11
scripts/monitor_estado/__init__.py
Normal file
@@ -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"
|
||||
122
scripts/monitor_estado/archivo.py
Normal file
122
scripts/monitor_estado/archivo.py
Normal file
@@ -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/<ruta>/<YYYY-MM-DD>.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 {}
|
||||
53
scripts/monitor_estado/aviso.py
Normal file
53
scripts/monitor_estado/aviso.py
Normal file
@@ -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
|
||||
133
scripts/monitor_estado/ciclo.py
Normal file
133
scripts/monitor_estado/ciclo.py
Normal file
@@ -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)
|
||||
19
scripts/monitor_estado/config.example.json
Normal file
19
scripts/monitor_estado/config.example.json
Normal file
@@ -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 }
|
||||
}
|
||||
117
scripts/monitor_estado/config.py
Normal file
117
scripts/monitor_estado/config.py
Normal file
@@ -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())
|
||||
2
scripts/monitor_estado/fuentes/__init__.py
Normal file
2
scripts/monitor_estado/fuentes/__init__.py
Normal file
@@ -0,0 +1,2 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Fuentes del monitor de estado (Telegram, correo, Jira)."""
|
||||
73
scripts/monitor_estado/fuentes/correo.py
Normal file
73
scripts/monitor_estado/fuentes/correo.py
Normal file
@@ -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
|
||||
94
scripts/monitor_estado/fuentes/jira.py
Normal file
94
scripts/monitor_estado/fuentes/jira.py
Normal file
@@ -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}"
|
||||
263
scripts/monitor_estado/fuentes/telegram.py
Normal file
263
scripts/monitor_estado/fuentes/telegram.py
Normal file
@@ -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)
|
||||
91
scripts/monitor_estado/lock.py
Normal file
91
scripts/monitor_estado/lock.py
Normal file
@@ -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))
|
||||
134
scripts/monitor_estado/resumen.py
Normal file
134
scripts/monitor_estado/resumen.py
Normal file
@@ -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": "<resumen en espanol, breve, indicando fuente y autor>", '
|
||||
'"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
|
||||
94
scripts/monitor_estado/vistos.py
Normal file
94
scripts/monitor_estado/vistos.py
Normal file
@@ -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/<clave>.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
|
||||
Reference in New Issue
Block a user