20-sep: migracion watchers (Telegram+recordatorios) a INFRAESTRUCTURA + estado/agentes/changelog

This commit is contained in:
juanminsanz
2026-09-20 12:42:16 +02:00
parent 5e8ae0c394
commit 143f9810b0
18 changed files with 1315 additions and 7 deletions

482
scripts/telegram_watcher.py Normal file
View File

@@ -0,0 +1,482 @@
"""
telegram_watcher.py — Muestra mensajes entrantes de Telegram en consola y gestiona canal IDD.
Avisos exclusivos al canal IDD:
- Aviso de conexión de Telegram al canal IDD
- Alertas de lectura (doble check) de contactos vigilados al canal IDD
- Solo contestar comandos dirigidos al canal IDD
- Soporte para silenciar/activar avisos (/silenciar, /activar)
"""
import asyncio, sys, os, json, time, re
from datetime import datetime, timezone
from pathlib import Path
C = {
'red': '\033[91m', 'green': '\033[92m', 'yellow': '\033[93m',
'blue': '\033[94m', 'magenta': '\033[95m', 'cyan': '\033[96m',
'white': '\033[97m', 'bold': '\033[1m', 'reset': '\033[0m',
'bg_red': '\033[41m',
}
SCRIPT_DIR = Path(__file__).resolve().parent
sys.path.insert(0, str(SCRIPT_DIR))
import infra_paths # noqa: E402 (resolucion de rutas del repo INFRAESTRUCTURA)
ESTADO_DIR = infra_paths.ESTADO_DIR
BIBLIOTECA_REPO = infra_paths.bootstrap_biblioteca_paths()
LOG_FILE = ESTADO_DIR / 'telegram_mensajes.log'
LAST_ID_FILE = ESTADO_DIR / '.watcher_last_id'
TRACKED_USERS_FILE = ESTADO_DIR / 'telegram_tracked_users.json'
ALERT_KEYWORDS = ['guardia', 'crítico', 'urgente', 'caído', 'error', 'sin asignar',
'RELEVANTE', 'incidencia', 'alarma', 'paro', 'parada']
ALARMA_COMMANDS = ['/alarma', '!alarma', '.alarma']
ALARMA_TWILIO_DIR = BIBLIOTECA_REPO / 'alarma_twilio'
SILENCIAR_COMMANDS = ['/silenciar', '!silenciar', '/mute', '!mute', '/pausar', '!pausar']
ACTIVAR_COMMANDS = ['/activar', '!activar', '/unmute', '!unmute', '/reanudar', '!reanudar']
NETATMO_COMMANDS = [
'/domotica', '!domotica', '/menu', '!menu',
'/netatmo', '!netatmo', '/ambiente', '!ambiente', '/clima', '!clima', '/sensores', '!sensores',
'/temperatura', '!temperatura', '/temp', '!temp',
'/humedad', '!humedad', '/hum', '!hum',
'/co2', '!co2',
'/ruido', '!ruido',
'/anunciar', '!anunciar',
'/alertas', '!alertas'
]
IDD_CHAT_ID = os.environ.get('TELEGRAM_CHAT_ID_2', '-1003967773087')
def load_tracked_config():
defaults = {
"contactos_vigilados": ["Romero", "Fernandez", "Dani", "Oscar", "Maribel"],
"avisar_canal_idd": True,
"silenciado": False,
"silenciado_hasta": 0,
"avisar_alexa": False,
"avisar_sonido_pc": False,
"cooldown_segundos": 60
}
if TRACKED_USERS_FILE.exists():
try:
with open(TRACKED_USERS_FILE, "r", encoding="utf-8") as f:
c = json.load(f)
defaults.update(c)
except Exception:
pass
return defaults
def save_tracked_config(c):
try:
with open(TRACKED_USERS_FILE, "w", encoding="utf-8") as f:
json.dump(c, f, indent=2)
except Exception as e:
print(f"Error guardando config tracked: {e}", file=sys.stderr)
def is_silenciado():
c = load_tracked_config()
if c.get("silenciado", False):
hasta = c.get("silenciado_hasta", 0)
if hasta == 0:
return True
if time.time() < hasta:
return True
else:
c["silenciado"] = False
c["silenciado_hasta"] = 0
save_tracked_config(c)
return False
return False
def play_beep():
try:
import winsound
winsound.Beep(1000, 200)
except:
pass
def send_to_idd(text, parse_mode="Markdown"):
"""Envia mensaje estrictamente al canal IDD."""
sys.path.insert(0, str(SCRIPT_DIR))
try:
import idd_channel
return idd_channel.send_message(text)
except Exception as e:
print(f"Error enviando a IDD: {e}", file=sys.stderr)
return None
def log_message(chat_name, sender_name, text, is_alert):
try:
entry = json.dumps({
'ts': datetime.now().isoformat(),
'chat': chat_name,
'sender': sender_name,
'text': text[:500] if text else '',
'alert': is_alert,
}, ensure_ascii=False)
with open(LOG_FILE, 'a', encoding='utf-8') as f:
f.write(entry + '\n')
except Exception as e:
print(f'{C["red"]}Error log: {e}{C["reset"]}')
async def handle_silenciar_command(text):
c = load_tracked_config()
parts = text.strip().split()
segundos = 0
tiempo_str = "de forma indefinida (hasta /activar)"
if len(parts) > 1:
param = parts[1].lower()
m_h = re.match(r'^(\d+)\s*h$', param)
m_m = re.match(r'^(\d+)\s*m$', param)
if m_h:
horas = int(m_h.group(1))
segundos = horas * 3600
tiempo_str = f"por {horas} hora(s)"
elif m_m:
minutos = int(m_m.group(1))
segundos = minutos * 60
tiempo_str = f"por {minutos} minuto(s)"
c["silenciado"] = True
c["silenciado_hasta"] = int(time.time() + segundos) if segundos > 0 else 0
save_tracked_config(c)
send_to_idd(f"🔕 *Avisos Silenciados*\nLos avisos de lectura y alertas en IDD han sido silenciados {tiempo_str}.\nUsa `/activar` para reanudar.")
async def handle_activar_command(text):
c = load_tracked_config()
c["silenciado"] = False
c["silenciado_hasta"] = 0
save_tracked_config(c)
send_to_idd("🔔 *Avisos Reactivados*\nLas notificaciones de lectura y alertas ambientales vuelven a estar activas en este canal.")
async def handle_alarma_command(text, chat_name, sender_name, client, dialog):
parts = text.strip().split(maxsplit=2)
if len(parts) < 2:
send_to_idd("ℹ️ Usos:\n!alarma +34657866417\n!alarma +34657866417 Mensaje personalizado")
return
numero = parts[1]
mensaje = parts[2] if len(parts) >= 3 else None
env_path = ALARMA_TWILIO_DIR / '.env'
if env_path.exists():
from dotenv import load_dotenv
load_dotenv(env_path)
sys.path.insert(0, str(ALARMA_TWILIO_DIR))
from config import validate_config
from twilio_service import TwilioAlarmaService
errores = validate_config()
if errores:
send_to_idd("❌ Twilio no configurado:\n" + "\n".join(errores))
return
msg = f"📞 Llamando a {numero}..."
if mensaje:
msg += f"\n📝 Mensaje: {mensaje[:80]}"
send_to_idd(msg)
def ejecutar():
svc = TwilioAlarmaService()
sid = svc.lanzar_llamada(numero, mensaje)
estado, contestado = svc.esperar_resultado(sid)
return estado, contestado
loop = asyncio.get_event_loop()
estado, contestado = await loop.run_in_executor(None, ejecutar)
emoji_map = {
"completed": "✅", "busy": "📞", "no-answer": "📵",
"failed": "❌", "canceled": "🚫", "timeout": "⏰",
}
emoji = emoji_map.get(estado, "❓")
desc = f"{emoji} {numero}: {estado}"
if contestado:
desc += f" ({contestado})"
if contestado == "human":
desc += " 🎯 Humano detectado"
elif contestado and "machine" in contestado:
desc += " 🤖 Contestador automático"
send_to_idd(desc)
async def handle_netatmo_command(text, chat_name, sender_name, client, dialog):
cmd = text.strip().lower().split()[0]
sys.path.insert(0, str(SCRIPT_DIR))
import netatmo_tool
def consultar():
try:
token = netatmo_tool.get_access_token()
data = netatmo_tool.get_homecoach_data(token)
sensores = netatmo_tool.extraer_sensores(data)
cfg = netatmo_tool.load_config()
return sensores, cfg, None
except Exception as e:
return None, None, str(e)
loop = asyncio.get_event_loop()
sensores, cfg, error = await loop.run_in_executor(None, consultar)
if error:
send_to_idd(f"⚠️ Error consultando Netatmo: {error}")
return
if cmd in ('/temperatura', '!temperatura', '/temp', '!temp'):
msg = netatmo_tool.generar_mensaje_telegram("temp", sensores, cfg)
elif cmd in ('/humedad', '!humedad', '/hum', '!hum'):
msg = netatmo_tool.generar_mensaje_telegram("hum", sensores, cfg)
elif cmd in ('/co2', '!co2'):
msg = netatmo_tool.generar_mensaje_telegram("co2", sensores, cfg)
elif cmd in ('/ruido', '!ruido'):
msg = netatmo_tool.generar_mensaje_telegram("ruido", sensores, cfg)
elif cmd in ('/anunciar', '!anunciar'):
netatmo_tool.cmd_anunciar(sensores)
msg = "📢 <b>Anuncio enviado a Alexa:</b>\nTemperatura y calidad ambiental leídas por altavoz Echo."
elif cmd in ('/alertas', '!alertas'):
alertas_voz, alertas_tg = netatmo_tool.evaluar_alertas(sensores, cfg)
if alertas_tg:
msg = "🚨 <b>PARÁMETROS FUERA DE RANGO:</b>\n" + "\n".join(f"• {a}" for a in alertas_tg)
else:
msg = "✅ <b>Todos los parámetros están dentro del rango seguro configurado.</b>"
else:
msg = netatmo_tool.generar_mensaje_telegram("menu", sensores, cfg)
netatmo_tool.notificar_telegram(msg)
async def main():
try:
from telethon import TelegramClient, events
except ImportError:
print(f'{C["red"]}pip install telethon{C["reset"]}')
sys.exit(1)
api_id = int(os.environ.get('TELEGRAM_API_ID', '0'))
api_hash = os.environ.get('TELEGRAM_API_HASH', '')
if not api_id or not api_hash:
import subprocess
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:
pass
if not api_id or not api_hash:
print(f'{C["red"]}Falta TELEGRAM_API_ID / TELEGRAM_API_HASH{C["reset"]}')
sys.exit(1)
session_file = str(Path.home() / '.telegram-agent' / 'sessions' / 'watcher')
client = TelegramClient(session_file, api_id, api_hash)
await client.start(phone='+34657866417')
me = await client.get_me()
print(f'{C["green"]}Conectado: {me.first_name} (@{me.username or "sin"}){C["reset"]}')
print(f'{C["cyan"]}Canal de avisos exclusivo: IDD ({IDD_CHAT_ID}){C["reset"]}')
print(f'{C["cyan"]}Log: {LOG_FILE}{C["reset"]}')
print(f'{"─" * 60}\n')
# Enviar aviso de conexión al canal IDD
send_to_idd(f"🟢 *Telegram Conectado* ({me.first_name})\nMonitorizando lecturas de contactos y canal IDD.")
# Recuperar ultimos IDs procesados por chat
last_ids = {}
if LAST_ID_FILE.exists():
try:
last_ids = json.loads(LAST_ID_FILE.read_text())
except:
pass
# Registro de eventos de confirmacion de lectura (Doble Check) -> SOLO IDD
last_read_triggers = {}
@client.on(events.MessageRead(outbox=True))
async def on_message_read(event):
try:
if is_silenciado():
return
c = load_tracked_config()
vigilados = c.get("contactos_vigilados", ["Romero", "Fernandez", "Dani", "Oscar", "Maribel"])
chat = await event.get_chat()
chat_title = getattr(chat, "title", "") or getattr(chat, "first_name", "") or str(event.chat_id)
username = getattr(chat, "username", "") or ""
coincide = False
nombre_vig = chat_title
for v in vigilados:
v_clean = v.lower().strip()
if v_clean and (v_clean in chat_title.lower() or (username and v_clean in username.lower()) or v_clean in str(event.chat_id)):
coincide = True
nombre_vig = v
break
if not coincide:
return
now_t = time.time()
if (now_t - last_read_triggers.get(event.chat_id, 0)) < c.get("cooldown_segundos", 60):
return
last_read_triggers[event.chat_id] = now_t
hora_now = datetime.now().strftime("%H:%M:%S")
print(f"\n{C['bg_red']}{C['white']} 👀 LEÍDO {C['reset']} {C['bold']}{C['green']}[{nombre_vig}]{C['reset']} {hora_now} ha leído tu mensaje en '{chat_title}'")
entry = json.dumps({
"ts": datetime.now().isoformat(),
"tipo": "LECTURA",
"contacto": nombre_vig,
"chat": chat_title
}, ensure_ascii=False)
with open(LOG_FILE, "a", encoding="utf-8") as f:
f.write(entry + "\n")
# Avisar exclusivamente al canal IDD
msg_idd = (
f"👀 *LECTURA CONFIRMADA*\n"
f"━━━━━━━━━━━━━━━━━━━━━━\n"
f"• *Contacto:* `{nombre_vig}`\n"
f"• *Chat:* {chat_title}\n"
f"• *Hora:* {hora_now}"
)
send_to_idd(msg_idd)
except Exception as e:
pass
print(f'{C["white"]}Escuchando... (Ctrl+C para salir){C["reset"]}\n')
last_netatmo_check = 0
while True:
try:
# Comprobacion periodica de alertas Netatmo (cada 5 minutos = 300s) -> SOLO IDD
now = time.time()
if now - last_netatmo_check > 300:
last_netatmo_check = now
try:
if not is_silenciado():
sys.path.insert(0, str(SCRIPT_DIR))
import netatmo_tool
def check_alertas():
tok = netatmo_tool.get_access_token()
dt = netatmo_tool.get_homecoach_data(tok)
sens = netatmo_tool.extraer_sensores(dt)
c = netatmo_tool.load_config()
return netatmo_tool.evaluar_alertas(sens, c), sens, c
loop = asyncio.get_event_loop()
(voz, tg), sens, cfg = await loop.run_in_executor(None, check_alertas)
if voz:
disparar = False
cooldown_sec = cfg.get("cooldown_minutos", 30) * 60
msg_voz = ". ".join(voz)
last_alert_file = ESTADO_DIR / ".netatmo_last_alert.json"
if last_alert_file.exists():
try:
with open(last_alert_file, "r", encoding="utf-8") as f:
last_d = json.load(f)
if last_d.get("msg") != msg_voz or (now - last_d.get("timestamp", 0)) >= cooldown_sec:
disparar = True
except Exception:
disparar = True
else:
disparar = True
if disparar:
# Notificar al canal IDD
alerta_html = (
f"🚨 <b>ALERTA DOMÓTICA — {sens[0].get('nombre', 'Netatmo')}</b>\n"
f"━━━━━━━━━━━━━━━━━━━━━━\n"
+ "\n".join(f"• {a}" for a in tg) + "\n"
f"━━━━━━━━━━━━━━━━━━━━━━\n"
f"<i>Detectado a las {datetime.now().strftime('%H:%M:%S')}</i>"
)
netatmo_tool.notificar_telegram(alerta_html)
with open(last_alert_file, "w", encoding="utf-8") as f:
json.dump({"timestamp": now, "msg": msg_voz, "fecha": datetime.now().isoformat()}, f, indent=2)
except Exception as e:
pass
dialogs = await client.get_dialogs(limit=30)
for dialog in dialogs:
try:
dialog_key = str(dialog.id)
last_id = last_ids.get(dialog_key, 0)
messages = await client.get_messages(dialog, limit=5)
new_max = last_id
for msg in reversed(messages):
if msg.id <= last_id or not msg.message:
continue
new_max = max(new_max, msg.id)
chat_name = dialog.name or str(dialog.id)
sender = await msg.get_sender()
sender_name = getattr(sender, 'first_name', '') or getattr(sender, 'title', '') or str(msg.sender_id) if sender else '??'
text = msg.message
is_alert = any(kw in (text or '').lower() for kw in ALERT_KEYWORDS)
alert_prefix = f'{C["bg_red"]}{C["white"]} ALERTA {C["reset"]} ' if is_alert else ''
chat_color = C['cyan'] if 'jira' in chat_name.lower() or 'informe' in chat_name.lower() else C['magenta'] if 'alerta' in chat_name.lower() else C['white']
ts = msg.date.strftime('%H:%M:%S') if msg.date else '--:--:--'
print(f'\n{alert_prefix}{C["bold"]}{chat_color}[{chat_name}]{C["reset"]} {ts} {C["yellow"]}{sender_name}{C["reset"]}')
for line in text.split('\n')[:15]:
print(f' {chat_color}│{C["reset"]} {line}')
if len(text.split('\n')) > 15:
print(f' {chat_color}│{C["reset"]} ...')
log_message(chat_name, sender_name, text, is_alert)
if is_alert:
play_beep()
# Solo procesar comandos si vienen del canal IDD o Saved Messages
is_idd_channel = (str(dialog.id) == str(IDD_CHAT_ID) or chat_name == "IDD" or "idd" in chat_name.lower())
if text and is_idd_channel:
clean_text = text.strip()
if any(clean_text.startswith(cmd) for cmd in SILENCIAR_COMMANDS):
await handle_silenciar_command(clean_text)
elif any(clean_text.startswith(cmd) for cmd in ACTIVAR_COMMANDS):
await handle_activar_command(clean_text)
elif any(clean_text.startswith(cmd) for cmd in ALARMA_COMMANDS):
await handle_alarma_command(text, chat_name, sender_name, client, dialog)
elif any(clean_text.startswith(cmd) for cmd in NETATMO_COMMANDS):
await handle_netatmo_command(text, chat_name, sender_name, client, dialog)
if new_max > last_id:
last_ids[dialog_key] = new_max
except Exception as e:
err = str(e)
if 'FLOOD_WAIT' in err:
wait = int(''.join(c for c in err if c.isdigit()) or '5')
print(f'{C["yellow"]}Rate limit {wait}s...{C["reset"]}')
await asyncio.sleep(wait)
elif 'AUTH_KEY' not in err and 'database' not in err.lower():
pass
LAST_ID_FILE.write_text(json.dumps(last_ids))
except Exception as e:
err = str(e)
if 'FLOOD' in err:
print(f'{C["yellow"]}Flood, esperando 10s...{C["reset"]}')
else:
print(f'{C["red"]}Error: {err[:100]}{C["reset"]}')
await asyncio.sleep(5)
if __name__ == '__main__':
try:
asyncio.run(main())
except KeyboardInterrupt:
print(f'\n{C["yellow"]}Cerrado.{C["reset"]}')