489 lines
20 KiB
Python
489 lines
20 KiB
Python
"""
|
||
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"]}')
|