Files
INFRAESTRUCTURA/scripts/telegram_watcher.py

489 lines
20 KiB
Python
Raw Blame History

This file contains invisible Unicode characters
This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
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
for _stream in (sys.stdout, sys.stderr):
try:
_stream.reconfigure(encoding='utf-8', errors='replace')
except Exception:
pass
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())
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"]}')