diff --git a/.env.example b/.env.example index c1129b0..4fe5002 100644 --- a/.env.example +++ b/.env.example @@ -4,7 +4,8 @@ BOT_TOKEN=123456:your-telegram-bot-token # Comma-separated Telegram user ids allowed to use the bot ALLOWED_IDS=111111111,222222222 -# Prometheus/Alertmanager endpoints (behind your own token-auth reverse-proxy path) +# Prometheus/Alertmanager endpoints behind your own reverse proxy. +# The bot sends "Authorization: Bearer " — the proxy must check it. AM_ALERTS_URL=https://your-domain/ambot/api/v2/alerts PROM_QUERY_URL=https://your-domain/prombot/api/v1/query AM_TOKEN=change-me @@ -15,3 +16,8 @@ SERVERS_ORDER=srv1,srv2,srv3 ALERT_POLL_SECONDS=45 ALERT_HTTP_TIMEOUT=15 +# Report "monitoring unreachable" after this many failed polls in a row +ALERT_WATCHDOG_FAILURES=4 + +# SQLite file with subscriptions and alert state +ALERTBOT_DB=notifications.db diff --git a/.gitattributes b/.gitattributes new file mode 100644 index 0000000..6313b56 --- /dev/null +++ b/.gitattributes @@ -0,0 +1 @@ +* text=auto eol=lf diff --git a/bot.py b/bot.py index 5f8ee72..a0e6997 100644 --- a/bot.py +++ b/bot.py @@ -1,51 +1,51 @@ -import asyncio -import logging -import sys - -from aiogram.types import BotCommand - -from init import bot, dp -from handlers import router, alert_poller - -logging.basicConfig( - level=logging.INFO, - format='%(asctime)s - %(name)s - %(levelname)s - %(message)s', - handlers=[logging.StreamHandler(sys.stdout), - logging.FileHandler('bot.log', encoding='utf-8')], -) -log = logging.getLogger(__name__) - -COMMANDS = [ - BotCommand(command="servers", description="CPU / RAM / диск / аптайм по узлам"), - BotCommand(command="services", description="что работает и что нет"), - BotCommand(command="certs", description="дни до истечения TLS"), - BotCommand(command="alerts", description="активные алерты (on/off — подписка)"), - BotCommand(command="status", description="что включено"), - BotCommand(command="help", description="справка"), -] - - -async def main() -> None: - log.info("alertbot starting") - dp.include_router(router) - await bot.delete_webhook(drop_pending_updates=True) - try: - await bot.set_my_commands(COMMANDS) - except Exception as e: - log.warning("set_my_commands failed: %s", e) - - poller = asyncio.create_task(alert_poller()) - try: - await dp.start_polling(bot, allowed_updates=dp.resolve_used_update_types(), - timeout=20, relax=0.1) - finally: - poller.cancel() - await bot.session.close() - log.info("alertbot stopped") - - -if __name__ == "__main__": - try: - asyncio.run(main()) - except (KeyboardInterrupt, SystemExit): - pass +import asyncio +import logging +import sys + +from aiogram.types import BotCommand + +from init import bot, dp +from handlers import router, alert_poller + +logging.basicConfig( + level=logging.INFO, + format='%(asctime)s - %(name)s - %(levelname)s - %(message)s', + handlers=[logging.StreamHandler(sys.stdout), + logging.FileHandler('bot.log', encoding='utf-8')], +) +log = logging.getLogger(__name__) + +COMMANDS = [ + BotCommand(command="servers", description="CPU / RAM / диск / аптайм по узлам"), + BotCommand(command="services", description="что работает и что нет"), + BotCommand(command="certs", description="дни до истечения TLS"), + BotCommand(command="alerts", description="активные алерты (on/off — подписка)"), + BotCommand(command="status", description="что включено"), + BotCommand(command="help", description="справка"), +] + + +async def main() -> None: + log.info("alertbot starting") + dp.include_router(router) + await bot.delete_webhook(drop_pending_updates=True) + try: + await bot.set_my_commands(COMMANDS) + except Exception as e: + log.warning("set_my_commands failed: %s", e) + + poller = asyncio.create_task(alert_poller()) + try: + await dp.start_polling(bot, allowed_updates=dp.resolve_used_update_types(), + timeout=20, relax=0.1) + finally: + poller.cancel() + await bot.session.close() + log.info("alertbot stopped") + + +if __name__ == "__main__": + try: + asyncio.run(main()) + except (KeyboardInterrupt, SystemExit): + pass diff --git a/handlers.py b/handlers.py index c29fab4..8aeb1cb 100644 --- a/handlers.py +++ b/handlers.py @@ -1,333 +1,376 @@ -import asyncio -import html -import logging -import os -import time -from typing import Dict, List - -import aiohttp -from aiogram import BaseMiddleware, Router -from aiogram.filters import Command, CommandObject -from aiogram.types import Message - -from init import (ALERT_HTTP_TIMEOUT, ALERT_POLL_SECONDS, ALLOWED_IDS, AM_ALERTS_URL, - AM_TOKEN, PROM_QUERY_URL, PROM_TOKEN, bot) -from sql import DatabaseManager - -router = Router() -db = DatabaseManager() -log = logging.getLogger(__name__) - -# Node names as used in Prometheus' "server" label — match your own scrape config. -SERVERS_ORDER = [s for s in os.environ.get("SERVERS_ORDER", "srv1,srv2,srv3").split(",") if s] - - -class AccessMiddleware(BaseMiddleware): - async def __call__(self, handler, event, data): - u = getattr(event, "from_user", None) or data.get("event_from_user") - if u is None or u.id in ALLOWED_IDS: - return await handler(event, data) - return None # silently ignore everyone else - - -router.message.outer_middleware(AccessMiddleware()) - - -# ====================================================================== -# small helpers -# ====================================================================== -def fmt_uptime(sec: float) -> str: - sec = int(max(0, sec)) - d, rem = divmod(sec, 86400) - h = rem // 3600 - return f"{d}д {h}ч" if d else f"{h}ч" - - -def settings_block(row) -> str: - _, _u, alerts_on = row - firing = len(db.list_firing_fingerprints()) - out = f"Алерты: {'включены' if alerts_on else 'выключены'}" - if firing: - out += f"\nСейчас активных алертов: {firing}" - return out - - -# ====================================================================== -# commands -# ====================================================================== -@router.message(Command("start")) -async def cmd_start(msg: Message): - uid = msg.from_user.id - db.add_user(uid, msg.from_user.username or msg.from_user.first_name or str(uid)) - row = db.get_user_info(uid) - await msg.answer( - "Мониторинг инфраструктуры\n\n" - "• сразу сообщает о проблемах и о возврате в норму (алерты);\n" - "• по запросу отдаёт состояние серверов, сервисов и сертификатов.\n\n" - + settings_block(row) + "\n\n" - "Команды — в меню слева от поля ввода, или /help." - ) - - -@router.message(Command("help")) -async def cmd_help(msg: Message): - await msg.answer( - "Команды\n\n" - "Сводки\n" - "/servers — CPU, RAM, диск, аптайм по узлам\n" - "/services — что работает и что нет\n" - "/certs — сколько дней осталось у TLS-сертификатов\n" - "/alerts — активные алерты сейчас\n\n" - "Настройки\n" - "/status — что включено\n" - "/alerts on | off — подписка на алерты" - ) - - -@router.message(Command("status")) -async def cmd_status(msg: Message): - row = db.get_user_info(msg.from_user.id) - if not row: - await msg.answer("Нажми /start.") - return - await msg.answer("Настройки\n\n" + settings_block(row)) - - -@router.message(Command("alerts")) -async def cmd_alerts(msg: Message, command: CommandObject): - uid = msg.from_user.id - if not db.get_user_info(uid): - db.add_user(uid, msg.from_user.username or str(uid)) - arg = (command.args or "").strip().lower() - if arg in ("on", "вкл"): - db.set_alerts(uid, True) - await msg.answer("Алерты включены.") - return - if arg in ("off", "выкл"): - db.set_alerts(uid, False) - await msg.answer("Алерты выключены.") - return - await msg.answer(await render_active_alerts()) - - -@router.message(Command("servers")) -async def cmd_servers(msg: Message): - await msg.answer(await render_servers()) - - -@router.message(Command("services")) -async def cmd_services(msg: Message): - await msg.answer(await render_services()) - - -@router.message(Command("certs")) -async def cmd_certs(msg: Message): - await msg.answer(await render_certs()) - - -# ====================================================================== -# Prometheus -# ====================================================================== -async def _promq(session: aiohttp.ClientSession, query: str) -> List[dict]: - async with session.get(PROM_QUERY_URL, params={"query": query, "token": PROM_TOKEN}, - timeout=aiohttp.ClientTimeout(total=ALERT_HTTP_TIMEOUT)) as r: - r.raise_for_status() - j = await r.json() - if j.get("status") != "success": - raise RuntimeError(j.get("error", "prometheus error")) - return j["data"]["result"] - - -def _by_server(result: List[dict]) -> Dict[str, float]: - out = {} - for s in result: - srv = s["metric"].get("server") - if srv: - try: - out[srv] = float(s["value"][1]) - except (ValueError, TypeError): - pass - return out - - -async def render_servers() -> str: - q = { - "cpu": '100 - (avg by (server) (rate(node_cpu_seconds_total{mode="idle"}[2m])) * 100)', - "ram": '(1 - node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes) * 100', - "disk": '(1 - avg by (server)(node_filesystem_avail_bytes{mountpoint="/",fstype!~"tmpfs|overlay|squashfs"}) ' - '/ avg by (server)(node_filesystem_size_bytes{mountpoint="/",fstype!~"tmpfs|overlay|squashfs"})) * 100', - "load": 'node_load1', - "up": 'up{job=~"node_.*"}', - "boot": 'node_boot_time_seconds', - "swap": '(node_memory_SwapTotal_bytes > bool 0) * (1 - node_memory_SwapFree_bytes / (node_memory_SwapTotal_bytes > 0)) * 100', - # real hwmon sensors only (chip=~"pci.*") — excludes bogus acpitz/thermal_zone - # readings some laptops-as-servers report as a stuck fake value - "temp": 'max by (server) (node_hwmon_temp_celsius{chip=~"pci.*"})', - } - try: - async with aiohttp.ClientSession() as s: - res = {k: _by_server(await _promq(s, v)) for k, v in q.items()} - except Exception as e: - return f"Не удалось получить метрики: {e}" - - now = time.time() - out = ["Серверы"] - for srv in SERVERS_ORDER: - if srv not in res["up"]: - out.append(f"\n{srv.upper()} — нет данных") - continue - alive = res["up"].get(srv, 0) >= 1 - cpu, ram = res["cpu"].get(srv), res["ram"].get(srv) - disk, swap = res["disk"].get(srv), res["swap"].get(srv, 0.0) - load = res["load"].get(srv) - temp = res["temp"].get(srv) - up = fmt_uptime(now - res["boot"][srv]) if srv in res["boot"] else "?" - - def g(x): - return f"{x:.0f}%" if isinstance(x, (int, float)) else "?" - - warn = [] - if isinstance(disk, float) and disk >= 85: - warn.append("диск") - if isinstance(ram, float) and ram >= 90: - warn.append("память") - if isinstance(swap, float) and swap >= 60: - warn.append("swap") - if isinstance(temp, float) and temp >= 85: - warn.append("температура") - head = f"\n{srv.upper()}" - if not alive: - head += " · НЕ ОТВЕЧАЕТ" - elif warn: - head += " · ⚠ " + ", ".join(warn) - out.append(head) - line = f"CPU {g(cpu)} · RAM {g(ram)} · диск {g(disk)} · swap {g(swap)}" - if isinstance(temp, float): - line += f" · temp {temp:.0f}°C" - out.append(line) - out.append(f"load {load:.2f} · аптайм {up}" if isinstance(load, float) - else f"аптайм {up}") - return "\n".join(out) - - -async def render_services() -> str: - try: - async with aiohttp.ClientSession() as s: - res = await _promq(s, "up") - except Exception as e: - return f"Не удалось получить статус сервисов: {e}" - total = len(res) - down = [(m["metric"].get("job", "?"), m["metric"].get("instance", "?")) - for m in res if m["value"][1] != "1"] - if not down: - return f"Сервисы\n\nВсе {total} в норме." - lines = [f"Сервисы\n\n{total - len(down)}/{total} в норме. Не отвечают:"] - for job, inst in sorted(down): - lines.append(f"• {html.escape(job)} — {html.escape(inst)}") - return "\n".join(lines) - - -async def render_certs() -> str: - try: - async with aiohttp.ClientSession() as s: - res = await _promq(s, "(probe_ssl_earliest_cert_expiry - time()) / 86400") - except Exception as e: - return f"Не удалось получить данные по сертификатам: {e}" - rows = [] - for m in res: - try: - days = float(m["value"][1]) - except (ValueError, TypeError): - continue - host = m["metric"].get("instance", "?").replace("https://", "").split("/")[0] - rows.append((days, host)) - if not rows: - return "Данных по сертификатам нет." - rows.sort() - soon = [f"{h} — {d:.0f} дн" for d, h in rows if d < 14] - body = "\n".join(f"{d:>4.0f} {h}" for d, h in rows) - head = "TLS-сертификаты" - if soon: - head += "\n⚠ скоро истекают: " + "; ".join(soon) - return head + f"\n
{body}
" - - -# ====================================================================== -# Alertmanager -# ====================================================================== -def _sev(labels): return labels.get("severity", "").lower() -def _where(labels): return labels.get("server") or labels.get("instance") or labels.get("job") or "" - - -def _fmt_firing(a: dict) -> str: - lb, an = a.get("labels", {}), a.get("annotations", {}) - name = html.escape(lb.get("alertname", "alert")) - sub = " · ".join(x for x in (_sev(lb), html.escape(_where(lb))) if x) - summ = html.escape(an.get("summary") or an.get("description") or "") - text = f"Проблема · {name}" - if sub: - text += f"\n{sub}" - if summ: - text += f"\n{summ}" - return text - - -def _fmt_resolved(name: str) -> str: - return f"В норме · {html.escape(name)}" - - -async def _fetch_alerts(session: aiohttp.ClientSession) -> List[dict]: - params = {"token": AM_TOKEN, "active": "true", "silenced": "false", "inhibited": "false"} - async with session.get(AM_ALERTS_URL, params=params, - timeout=aiohttp.ClientTimeout(total=ALERT_HTTP_TIMEOUT)) as r: - r.raise_for_status() - return await r.json() - - -async def _broadcast(text: str): - for uid in db.get_alert_users(): - try: - await bot.send_message(uid, text) - except Exception as e: - log.error("alert to %s failed: %s", uid, e) - - -async def render_active_alerts() -> str: - try: - async with aiohttp.ClientSession() as s: - data = await _fetch_alerts(s) - except Exception as e: - return f"Не удалось получить алерты: {e}" - firing = [a for a in data if a.get("status", {}).get("state") == "active"] - if not firing: - return "Активных алертов нет." - return f"Активные алерты: {len(firing)}\n\n" + "\n\n".join(_fmt_firing(a) for a in firing) - - -async def alert_poller(): - log.info("alert poller started (every %ss)", ALERT_POLL_SECONDS) - async with aiohttp.ClientSession() as session: - while True: - try: - data = await _fetch_alerts(session) - current = {} - for a in data: - if a.get("status", {}).get("state") != "active": - continue - fp = a.get("fingerprint") - if not fp: - continue - current[fp] = a - if db.get_alert_status(fp) != "firing": - await _broadcast(_fmt_firing(a)) - db.upsert_alert(fp, "firing", a.get("labels", {}).get("alertname", "alert")) - for fp, name in db.list_firing_fingerprints(): - if fp not in current: - await _broadcast(_fmt_resolved(name)) - db.upsert_alert(fp, "resolved", name) - db.purge_old_resolved() - except asyncio.CancelledError: - break - except Exception as e: - log.warning("alert poll failed: %s", e) - await asyncio.sleep(ALERT_POLL_SECONDS) +import asyncio +import html +import logging +import os +import time +from typing import Dict, List + +import aiohttp +from aiogram import BaseMiddleware, Router +from aiogram.filters import Command, CommandObject +from aiogram.types import Message + +from init import (ALERT_HTTP_TIMEOUT, ALERT_POLL_SECONDS, ALERT_WATCHDOG_FAILURES, ALLOWED_IDS, + AM_ALERTS_URL, AM_TOKEN, DB_PATH, PROM_QUERY_URL, PROM_TOKEN, bot) +from sql import DatabaseManager + +router = Router() +db = DatabaseManager(DB_PATH) +log = logging.getLogger(__name__) + +# Node names as used in Prometheus' "server" label — match your own scrape config. +SERVERS_ORDER = [s for s in os.environ.get("SERVERS_ORDER", "srv1,srv2,srv3").split(",") if s] + + +class AccessMiddleware(BaseMiddleware): + async def __call__(self, handler, event, data): + u = getattr(event, "from_user", None) or data.get("event_from_user") + if u is None or u.id in ALLOWED_IDS: + return await handler(event, data) + return None # silently ignore everyone else + + +router.message.outer_middleware(AccessMiddleware()) + + +# ====================================================================== +# small helpers +# ====================================================================== +def fmt_uptime(sec: float) -> str: + sec = int(max(0, sec)) + d, rem = divmod(sec, 86400) + h = rem // 3600 + return f"{d}д {h}ч" if d else f"{h}ч" + + +def settings_block(row) -> str: + _, _u, alerts_on = row + firing = len(db.list_firing_fingerprints()) + out = f"Алерты: {'включены' if alerts_on else 'выключены'}" + if firing: + out += f"\nСейчас активных алертов: {firing}" + return out + + +# ====================================================================== +# commands +# ====================================================================== +@router.message(Command("start")) +async def cmd_start(msg: Message): + uid = msg.from_user.id + db.add_user(uid, msg.from_user.username or msg.from_user.first_name or str(uid)) + row = db.get_user_info(uid) + await msg.answer( + "Мониторинг инфраструктуры\n\n" + "• сразу сообщает о проблемах и о возврате в норму (алерты);\n" + "• по запросу отдаёт состояние серверов, сервисов и сертификатов.\n\n" + + settings_block(row) + "\n\n" + "Команды — в меню слева от поля ввода, или /help." + ) + + +@router.message(Command("help")) +async def cmd_help(msg: Message): + await msg.answer( + "Команды\n\n" + "Сводки\n" + "/servers — CPU, RAM, диск, аптайм по узлам\n" + "/services — что работает и что нет\n" + "/certs — сколько дней осталось у TLS-сертификатов\n" + "/alerts — активные алерты сейчас\n\n" + "Настройки\n" + "/status — что включено\n" + "/alerts on | off — подписка на алерты" + ) + + +@router.message(Command("status")) +async def cmd_status(msg: Message): + row = db.get_user_info(msg.from_user.id) + if not row: + await msg.answer("Нажми /start.") + return + await msg.answer("Настройки\n\n" + settings_block(row)) + + +@router.message(Command("alerts")) +async def cmd_alerts(msg: Message, command: CommandObject): + uid = msg.from_user.id + if not db.get_user_info(uid): + db.add_user(uid, msg.from_user.username or str(uid)) + arg = (command.args or "").strip().lower() + if arg in ("on", "вкл"): + db.set_alerts(uid, True) + await msg.answer("Алерты включены.") + return + if arg in ("off", "выкл"): + db.set_alerts(uid, False) + await msg.answer("Алерты выключены.") + return + await msg.answer(await render_active_alerts()) + + +@router.message(Command("servers")) +async def cmd_servers(msg: Message): + await msg.answer(await render_servers()) + + +@router.message(Command("services")) +async def cmd_services(msg: Message): + await msg.answer(await render_services()) + + +@router.message(Command("certs")) +async def cmd_certs(msg: Message): + await msg.answer(await render_certs()) + + +# ====================================================================== +# Prometheus +# ====================================================================== +def _auth(token: str) -> dict: + """Token goes in the Authorization header, not in the URL (URLs end up in proxy logs).""" + return {"Authorization": f"Bearer {token}"} + + +async def _promq(session: aiohttp.ClientSession, query: str) -> List[dict]: + async with session.get(PROM_QUERY_URL, params={"query": query}, + headers=_auth(PROM_TOKEN), + timeout=aiohttp.ClientTimeout(total=ALERT_HTTP_TIMEOUT)) as r: + r.raise_for_status() + j = await r.json() + if j.get("status") != "success": + raise RuntimeError(j.get("error", "prometheus error")) + return j["data"]["result"] + + +def _by_server(result: List[dict]) -> Dict[str, float]: + out = {} + for s in result: + srv = s["metric"].get("server") + if srv: + try: + out[srv] = float(s["value"][1]) + except (ValueError, TypeError): + pass + return out + + +async def render_servers() -> str: + q = { + "cpu": '100 - (avg by (server) (rate(node_cpu_seconds_total{mode="idle"}[2m])) * 100)', + "ram": '(1 - node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes) * 100', + "disk": '(1 - avg by (server)(node_filesystem_avail_bytes{mountpoint="/",fstype!~"tmpfs|overlay|squashfs"}) ' + '/ avg by (server)(node_filesystem_size_bytes{mountpoint="/",fstype!~"tmpfs|overlay|squashfs"})) * 100', + "load": 'node_load1', + "up": 'up{job=~"node_.*"}', + "boot": 'node_boot_time_seconds', + "swap": '(node_memory_SwapTotal_bytes > bool 0) * (1 - node_memory_SwapFree_bytes / (node_memory_SwapTotal_bytes > 0)) * 100', + # real hwmon sensors only (chip=~"pci.*") — excludes bogus acpitz/thermal_zone + # readings some laptops-as-servers report as a stuck fake value + "temp": 'max by (server) (node_hwmon_temp_celsius{chip=~"pci.*"})', + } + try: + async with aiohttp.ClientSession() as s: + res = {k: _by_server(await _promq(s, v)) for k, v in q.items()} + except Exception as e: + return f"Не удалось получить метрики: {e}" + + now = time.time() + out = ["Серверы"] + for srv in SERVERS_ORDER: + if srv not in res["up"]: + out.append(f"\n{srv.upper()} — нет данных") + continue + alive = res["up"].get(srv, 0) >= 1 + cpu, ram = res["cpu"].get(srv), res["ram"].get(srv) + disk, swap = res["disk"].get(srv), res["swap"].get(srv, 0.0) + load = res["load"].get(srv) + temp = res["temp"].get(srv) + up = fmt_uptime(now - res["boot"][srv]) if srv in res["boot"] else "?" + + def g(x): + return f"{x:.0f}%" if isinstance(x, (int, float)) else "?" + + warn = [] + if isinstance(disk, float) and disk >= 85: + warn.append("диск") + if isinstance(ram, float) and ram >= 90: + warn.append("память") + if isinstance(swap, float) and swap >= 60: + warn.append("swap") + if isinstance(temp, float) and temp >= 85: + warn.append("температура") + head = f"\n{srv.upper()}" + if not alive: + head += " · НЕ ОТВЕЧАЕТ" + elif warn: + head += " · ⚠ " + ", ".join(warn) + out.append(head) + line = f"CPU {g(cpu)} · RAM {g(ram)} · диск {g(disk)} · swap {g(swap)}" + if isinstance(temp, float): + line += f" · temp {temp:.0f}°C" + out.append(line) + out.append(f"load {load:.2f} · аптайм {up}" if isinstance(load, float) + else f"аптайм {up}") + return "\n".join(out) + + +async def render_services() -> str: + try: + async with aiohttp.ClientSession() as s: + res = await _promq(s, "up") + except Exception as e: + return f"Не удалось получить статус сервисов: {e}" + total = len(res) + down = [(m["metric"].get("job", "?"), m["metric"].get("instance", "?")) + for m in res if m["value"][1] != "1"] + if not down: + return f"Сервисы\n\nВсе {total} в норме." + lines = [f"Сервисы\n\n{total - len(down)}/{total} в норме. Не отвечают:"] + for job, inst in sorted(down): + lines.append(f"• {html.escape(job)} — {html.escape(inst)}") + return "\n".join(lines) + + +async def render_certs() -> str: + try: + async with aiohttp.ClientSession() as s: + res = await _promq(s, "(probe_ssl_earliest_cert_expiry - time()) / 86400") + except Exception as e: + return f"Не удалось получить данные по сертификатам: {e}" + rows = [] + for m in res: + try: + days = float(m["value"][1]) + except (ValueError, TypeError): + continue + host = m["metric"].get("instance", "?").replace("https://", "").split("/")[0] + rows.append((days, host)) + if not rows: + return "Данных по сертификатам нет." + rows.sort() + soon = [f"{h} — {d:.0f} дн" for d, h in rows if d < 14] + body = "\n".join(f"{d:>4.0f} {h}" for d, h in rows) + head = "TLS-сертификаты" + if soon: + head += "\n⚠ скоро истекают: " + "; ".join(soon) + return head + f"\n
{body}
" + + +# ====================================================================== +# Alertmanager +# ====================================================================== +def _sev(labels): return labels.get("severity", "").lower() +def _where(labels): return labels.get("server") or labels.get("instance") or labels.get("job") or "" + + +def _fmt_firing(a: dict) -> str: + lb, an = a.get("labels", {}), a.get("annotations", {}) + name = html.escape(lb.get("alertname", "alert")) + sub = " · ".join(x for x in (_sev(lb), html.escape(_where(lb))) if x) + summ = html.escape(an.get("summary") or an.get("description") or "") + text = f"Проблема · {name}" + if sub: + text += f"\n{sub}" + if summ: + text += f"\n{summ}" + return text + + +def _fmt_resolved(name: str) -> str: + return f"В норме · {html.escape(name)}" + + +async def _fetch_alerts(session: aiohttp.ClientSession) -> List[dict]: + params = {"active": "true", "silenced": "false", "inhibited": "false"} + async with session.get(AM_ALERTS_URL, params=params, headers=_auth(AM_TOKEN), + timeout=aiohttp.ClientTimeout(total=ALERT_HTTP_TIMEOUT)) as r: + r.raise_for_status() + return await r.json() + + +async def _broadcast(text: str): + for uid in db.get_alert_users(): + try: + await bot.send_message(uid, text) + except Exception as e: + log.error("alert to %s failed: %s", uid, e) + + +async def render_active_alerts() -> str: + try: + async with aiohttp.ClientSession() as s: + data = await _fetch_alerts(s) + except Exception as e: + return f"Не удалось получить алерты: {e}" + firing = [a for a in data if a.get("status", {}).get("state") == "active"] + if not firing: + return "Активных алертов нет." + return f"Активные алерты: {len(firing)}\n\n" + "\n\n".join(_fmt_firing(a) for a in firing) + + +def _fmt_watchdog_down(error: str) -> str: + return ("Мониторинг недоступен\nAlertmanager не отвечает, о новых проблемах бот сейчас не узнает.\n" + f"Последняя ошибка: {html.escape(error)}") + + +def _fmt_watchdog_up() -> str: + return "Мониторинг снова доступен" + + +async def process_alerts(data: List[dict]) -> None: + """One poll: announce new firing alerts and alerts that have resolved since the last poll.""" + current = {} + for a in data: + if a.get("status", {}).get("state") != "active": + continue + fp = a.get("fingerprint") + if not fp: + continue + current[fp] = a + if db.get_alert_status(fp) != "firing": + await _broadcast(_fmt_firing(a)) + db.upsert_alert(fp, "firing", a.get("labels", {}).get("alertname", "alert")) + for fp, name in db.list_firing_fingerprints(): + if fp not in current: + await _broadcast(_fmt_resolved(name)) + db.upsert_alert(fp, "resolved", name) + db.purge_old_resolved() + + +class Watchdog: + """Counts failed polls in a row; reports once when monitoring is lost and once when it is back.""" + + def __init__(self, threshold: int): + self.threshold = threshold + self.failures = 0 + self.reported = False + + async def failed(self, error: str) -> None: + self.failures += 1 + if self.failures >= self.threshold and not self.reported: + self.reported = True + await _broadcast(_fmt_watchdog_down(error)) + + async def ok(self) -> None: + if self.reported: + await _broadcast(_fmt_watchdog_up()) + self.failures = 0 + self.reported = False + + +async def alert_poller(): + log.info("alert poller started (every %ss)", ALERT_POLL_SECONDS) + watchdog = Watchdog(ALERT_WATCHDOG_FAILURES) + async with aiohttp.ClientSession() as session: + while True: + try: + await process_alerts(await _fetch_alerts(session)) + await watchdog.ok() + except asyncio.CancelledError: + break + except Exception as e: + log.warning("alert poll failed: %s", e) + await watchdog.failed(str(e) or type(e).__name__) + await asyncio.sleep(ALERT_POLL_SECONDS) diff --git a/init.py b/init.py index bac6eb8..7b20fa9 100644 --- a/init.py +++ b/init.py @@ -20,6 +20,9 @@ AM_TOKEN = os.environ["AM_TOKEN"] PROM_TOKEN = os.environ["PROM_TOKEN"] ALERT_POLL_SECONDS = int(os.environ.get("ALERT_POLL_SECONDS", "45")) ALERT_HTTP_TIMEOUT = int(os.environ.get("ALERT_HTTP_TIMEOUT", "15")) +# After this many failed polls in a row the bot reports that monitoring itself is unreachable +ALERT_WATCHDOG_FAILURES = int(os.environ.get("ALERT_WATCHDOG_FAILURES", "4")) +DB_PATH = os.environ.get("ALERTBOT_DB", "notifications.db") bot = Bot(token=BOT_TOKEN, default=DefaultBotProperties(parse_mode=ParseMode.HTML)) dp = Dispatcher(storage=MemoryStorage()) diff --git a/sql.py b/sql.py index a8d06e8..550aa5e 100644 --- a/sql.py +++ b/sql.py @@ -1,118 +1,118 @@ -import sqlite3 -from typing import List, Optional, Tuple - - -class DatabaseManager: - def __init__(self, db_path: str = "notifications.db"): - self.db_path = db_path - self.init_database() - - def _conn(self): - return sqlite3.connect(self.db_path) - - def init_database(self): - with self._conn() as conn: - conn.execute(''' - CREATE TABLE IF NOT EXISTS users ( - user_id INTEGER PRIMARY KEY, - username TEXT, - alerts_enabled BOOLEAN DEFAULT 1, - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - ''') - conn.execute(''' - CREATE TABLE IF NOT EXISTS alert_seen ( - fingerprint TEXT PRIMARY KEY, - status TEXT, - name TEXT, - updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - ''') - conn.commit() - - # ---------- users ---------- - def add_user(self, user_id: int, username: str) -> bool: - try: - with self._conn() as conn: - conn.execute( - "INSERT INTO users (user_id, username) VALUES (?, ?) " - "ON CONFLICT(user_id) DO UPDATE SET username=excluded.username", - (user_id, username), - ) - conn.commit() - return True - except sqlite3.Error: - return False - - def set_alerts(self, user_id: int, value: bool) -> bool: - try: - with self._conn() as conn: - conn.execute("UPDATE users SET alerts_enabled = ? WHERE user_id = ?", - (1 if value else 0, user_id)) - conn.commit() - return True - except sqlite3.Error: - return False - - def get_user_info(self, user_id: int) -> Optional[Tuple]: - """(user_id, username, alerts_enabled)""" - try: - with self._conn() as conn: - return conn.execute( - "SELECT user_id, username, alerts_enabled FROM users WHERE user_id = ?", - (user_id,) - ).fetchone() - except sqlite3.Error: - return None - - def get_alert_users(self) -> List[int]: - try: - with self._conn() as conn: - return [r[0] for r in conn.execute( - "SELECT user_id FROM users WHERE alerts_enabled = 1" - ).fetchall()] - except sqlite3.Error: - return [] - - # ---------- alert dedup ---------- - def get_alert_status(self, fingerprint: str) -> Optional[str]: - try: - with self._conn() as conn: - r = conn.execute("SELECT status FROM alert_seen WHERE fingerprint = ?", - (fingerprint,)).fetchone() - return r[0] if r else None - except sqlite3.Error: - return None - - def upsert_alert(self, fingerprint: str, status: str, name: str): - try: - with self._conn() as conn: - conn.execute( - "INSERT INTO alert_seen (fingerprint, status, name, updated_at) " - "VALUES (?, ?, ?, CURRENT_TIMESTAMP) " - "ON CONFLICT(fingerprint) DO UPDATE SET status=excluded.status, " - "name=excluded.name, updated_at=CURRENT_TIMESTAMP", - (fingerprint, status, name), - ) - conn.commit() - except sqlite3.Error: - pass - - def list_firing_fingerprints(self) -> List[Tuple[str, str]]: - try: - with self._conn() as conn: - return conn.execute( - "SELECT fingerprint, name FROM alert_seen WHERE status = 'firing'" - ).fetchall() - except sqlite3.Error: - return [] - - def purge_old_resolved(self, days: int = 3): - try: - with self._conn() as conn: - conn.execute( - "DELETE FROM alert_seen WHERE status = 'resolved' " - "AND updated_at < datetime('now', ?)", (f'-{days} days',)) - conn.commit() - except sqlite3.Error: - pass +import sqlite3 +from typing import List, Optional, Tuple + + +class DatabaseManager: + def __init__(self, db_path: str = "notifications.db"): + self.db_path = db_path + self.init_database() + + def _conn(self): + return sqlite3.connect(self.db_path) + + def init_database(self): + with self._conn() as conn: + conn.execute(''' + CREATE TABLE IF NOT EXISTS users ( + user_id INTEGER PRIMARY KEY, + username TEXT, + alerts_enabled BOOLEAN DEFAULT 1, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP + ) + ''') + conn.execute(''' + CREATE TABLE IF NOT EXISTS alert_seen ( + fingerprint TEXT PRIMARY KEY, + status TEXT, + name TEXT, + updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP + ) + ''') + conn.commit() + + # ---------- users ---------- + def add_user(self, user_id: int, username: str) -> bool: + try: + with self._conn() as conn: + conn.execute( + "INSERT INTO users (user_id, username) VALUES (?, ?) " + "ON CONFLICT(user_id) DO UPDATE SET username=excluded.username", + (user_id, username), + ) + conn.commit() + return True + except sqlite3.Error: + return False + + def set_alerts(self, user_id: int, value: bool) -> bool: + try: + with self._conn() as conn: + conn.execute("UPDATE users SET alerts_enabled = ? WHERE user_id = ?", + (1 if value else 0, user_id)) + conn.commit() + return True + except sqlite3.Error: + return False + + def get_user_info(self, user_id: int) -> Optional[Tuple]: + """(user_id, username, alerts_enabled)""" + try: + with self._conn() as conn: + return conn.execute( + "SELECT user_id, username, alerts_enabled FROM users WHERE user_id = ?", + (user_id,) + ).fetchone() + except sqlite3.Error: + return None + + def get_alert_users(self) -> List[int]: + try: + with self._conn() as conn: + return [r[0] for r in conn.execute( + "SELECT user_id FROM users WHERE alerts_enabled = 1" + ).fetchall()] + except sqlite3.Error: + return [] + + # ---------- alert dedup ---------- + def get_alert_status(self, fingerprint: str) -> Optional[str]: + try: + with self._conn() as conn: + r = conn.execute("SELECT status FROM alert_seen WHERE fingerprint = ?", + (fingerprint,)).fetchone() + return r[0] if r else None + except sqlite3.Error: + return None + + def upsert_alert(self, fingerprint: str, status: str, name: str): + try: + with self._conn() as conn: + conn.execute( + "INSERT INTO alert_seen (fingerprint, status, name, updated_at) " + "VALUES (?, ?, ?, CURRENT_TIMESTAMP) " + "ON CONFLICT(fingerprint) DO UPDATE SET status=excluded.status, " + "name=excluded.name, updated_at=CURRENT_TIMESTAMP", + (fingerprint, status, name), + ) + conn.commit() + except sqlite3.Error: + pass + + def list_firing_fingerprints(self) -> List[Tuple[str, str]]: + try: + with self._conn() as conn: + return conn.execute( + "SELECT fingerprint, name FROM alert_seen WHERE status = 'firing'" + ).fetchall() + except sqlite3.Error: + return [] + + def purge_old_resolved(self, days: int = 3): + try: + with self._conn() as conn: + conn.execute( + "DELETE FROM alert_seen WHERE status = 'resolved' " + "AND updated_at < datetime('now', ?)", (f'-{days} days',)) + conn.commit() + except sqlite3.Error: + pass