1
0
Fork 0
mirror of https://github.com/EDeev/alertbot.git synced 2026-10-07 20:49:56 +03:00

Токены в заголовке, сторож недоступности мониторинга, путь к базе из окружения

- токены Prometheus/Alertmanager передаются в заголовке Authorization, а не в адресе;
- если Alertmanager не отвечает ALERT_WATCHDOG_FAILURES опросов подряд, бот сообщает
  «Мониторинг недоступен» и отдельно — когда он снова доступен;
- один цикл опроса вынесен в process_alerts(); путь к SQLite — ALERTBOT_DB;
- переводы строк приведены к LF.
This commit is contained in:
Egor Deev 2026-10-02 13:54:49 +00:00
parent 1a77cf8620
commit 63c9976888
6 changed files with 556 additions and 503 deletions

View file

@ -4,7 +4,8 @@ BOT_TOKEN=123456:your-telegram-bot-token
# Comma-separated Telegram user ids allowed to use the bot # Comma-separated Telegram user ids allowed to use the bot
ALLOWED_IDS=111111111,222222222 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 <token>" — the proxy must check it.
AM_ALERTS_URL=https://your-domain/ambot/api/v2/alerts AM_ALERTS_URL=https://your-domain/ambot/api/v2/alerts
PROM_QUERY_URL=https://your-domain/prombot/api/v1/query PROM_QUERY_URL=https://your-domain/prombot/api/v1/query
AM_TOKEN=change-me AM_TOKEN=change-me
@ -15,3 +16,8 @@ SERVERS_ORDER=srv1,srv2,srv3
ALERT_POLL_SECONDS=45 ALERT_POLL_SECONDS=45
ALERT_HTTP_TIMEOUT=15 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

1
.gitattributes vendored Normal file
View file

@ -0,0 +1 @@
* text=auto eol=lf

102
bot.py
View file

@ -1,51 +1,51 @@
import asyncio import asyncio
import logging import logging
import sys import sys
from aiogram.types import BotCommand from aiogram.types import BotCommand
from init import bot, dp from init import bot, dp
from handlers import router, alert_poller from handlers import router, alert_poller
logging.basicConfig( logging.basicConfig(
level=logging.INFO, level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s', format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
handlers=[logging.StreamHandler(sys.stdout), handlers=[logging.StreamHandler(sys.stdout),
logging.FileHandler('bot.log', encoding='utf-8')], logging.FileHandler('bot.log', encoding='utf-8')],
) )
log = logging.getLogger(__name__) log = logging.getLogger(__name__)
COMMANDS = [ COMMANDS = [
BotCommand(command="servers", description="CPU / RAM / диск / аптайм по узлам"), BotCommand(command="servers", description="CPU / RAM / диск / аптайм по узлам"),
BotCommand(command="services", description="что работает и что нет"), BotCommand(command="services", description="что работает и что нет"),
BotCommand(command="certs", description="дни до истечения TLS"), BotCommand(command="certs", description="дни до истечения TLS"),
BotCommand(command="alerts", description="активные алерты (on/off — подписка)"), BotCommand(command="alerts", description="активные алерты (on/off — подписка)"),
BotCommand(command="status", description="что включено"), BotCommand(command="status", description="что включено"),
BotCommand(command="help", description="справка"), BotCommand(command="help", description="справка"),
] ]
async def main() -> None: async def main() -> None:
log.info("alertbot starting") log.info("alertbot starting")
dp.include_router(router) dp.include_router(router)
await bot.delete_webhook(drop_pending_updates=True) await bot.delete_webhook(drop_pending_updates=True)
try: try:
await bot.set_my_commands(COMMANDS) await bot.set_my_commands(COMMANDS)
except Exception as e: except Exception as e:
log.warning("set_my_commands failed: %s", e) log.warning("set_my_commands failed: %s", e)
poller = asyncio.create_task(alert_poller()) poller = asyncio.create_task(alert_poller())
try: try:
await dp.start_polling(bot, allowed_updates=dp.resolve_used_update_types(), await dp.start_polling(bot, allowed_updates=dp.resolve_used_update_types(),
timeout=20, relax=0.1) timeout=20, relax=0.1)
finally: finally:
poller.cancel() poller.cancel()
await bot.session.close() await bot.session.close()
log.info("alertbot stopped") log.info("alertbot stopped")
if __name__ == "__main__": if __name__ == "__main__":
try: try:
asyncio.run(main()) asyncio.run(main())
except (KeyboardInterrupt, SystemExit): except (KeyboardInterrupt, SystemExit):
pass pass

View file

@ -1,333 +1,376 @@
import asyncio import asyncio
import html import html
import logging import logging
import os import os
import time import time
from typing import Dict, List from typing import Dict, List
import aiohttp import aiohttp
from aiogram import BaseMiddleware, Router from aiogram import BaseMiddleware, Router
from aiogram.filters import Command, CommandObject from aiogram.filters import Command, CommandObject
from aiogram.types import Message from aiogram.types import Message
from init import (ALERT_HTTP_TIMEOUT, ALERT_POLL_SECONDS, ALLOWED_IDS, AM_ALERTS_URL, from init import (ALERT_HTTP_TIMEOUT, ALERT_POLL_SECONDS, ALERT_WATCHDOG_FAILURES, ALLOWED_IDS,
AM_TOKEN, PROM_QUERY_URL, PROM_TOKEN, bot) AM_ALERTS_URL, AM_TOKEN, DB_PATH, PROM_QUERY_URL, PROM_TOKEN, bot)
from sql import DatabaseManager from sql import DatabaseManager
router = Router() router = Router()
db = DatabaseManager() db = DatabaseManager(DB_PATH)
log = logging.getLogger(__name__) log = logging.getLogger(__name__)
# Node names as used in Prometheus' "server" label — match your own scrape config. # 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] SERVERS_ORDER = [s for s in os.environ.get("SERVERS_ORDER", "srv1,srv2,srv3").split(",") if s]
class AccessMiddleware(BaseMiddleware): class AccessMiddleware(BaseMiddleware):
async def __call__(self, handler, event, data): async def __call__(self, handler, event, data):
u = getattr(event, "from_user", None) or data.get("event_from_user") u = getattr(event, "from_user", None) or data.get("event_from_user")
if u is None or u.id in ALLOWED_IDS: if u is None or u.id in ALLOWED_IDS:
return await handler(event, data) return await handler(event, data)
return None # silently ignore everyone else return None # silently ignore everyone else
router.message.outer_middleware(AccessMiddleware()) router.message.outer_middleware(AccessMiddleware())
# ====================================================================== # ======================================================================
# small helpers # small helpers
# ====================================================================== # ======================================================================
def fmt_uptime(sec: float) -> str: def fmt_uptime(sec: float) -> str:
sec = int(max(0, sec)) sec = int(max(0, sec))
d, rem = divmod(sec, 86400) d, rem = divmod(sec, 86400)
h = rem // 3600 h = rem // 3600
return f"{d}д {h}ч" if d else f"{h}ч" return f"{d}д {h}ч" if d else f"{h}ч"
def settings_block(row) -> str: def settings_block(row) -> str:
_, _u, alerts_on = row _, _u, alerts_on = row
firing = len(db.list_firing_fingerprints()) firing = len(db.list_firing_fingerprints())
out = f"Алерты: {'включены' if alerts_on else 'выключены'}" out = f"Алерты: {'включены' if alerts_on else 'выключены'}"
if firing: if firing:
out += f"\nСейчас активных алертов: <b>{firing}</b>" out += f"\nСейчас активных алертов: <b>{firing}</b>"
return out return out
# ====================================================================== # ======================================================================
# commands # commands
# ====================================================================== # ======================================================================
@router.message(Command("start")) @router.message(Command("start"))
async def cmd_start(msg: Message): async def cmd_start(msg: Message):
uid = msg.from_user.id uid = msg.from_user.id
db.add_user(uid, msg.from_user.username or msg.from_user.first_name or str(uid)) db.add_user(uid, msg.from_user.username or msg.from_user.first_name or str(uid))
row = db.get_user_info(uid) row = db.get_user_info(uid)
await msg.answer( await msg.answer(
"<b>Мониторинг инфраструктуры</b>\n\n" "<b>Мониторинг инфраструктуры</b>\n\n"
"• сразу сообщает о проблемах и о возврате в норму (алерты);\n" "• сразу сообщает о проблемах и о возврате в норму (алерты);\n"
"• по запросу отдаёт состояние серверов, сервисов и сертификатов.\n\n" "• по запросу отдаёт состояние серверов, сервисов и сертификатов.\n\n"
+ settings_block(row) + "\n\n" + settings_block(row) + "\n\n"
"Команды — в меню слева от поля ввода, или /help." "Команды — в меню слева от поля ввода, или /help."
) )
@router.message(Command("help")) @router.message(Command("help"))
async def cmd_help(msg: Message): async def cmd_help(msg: Message):
await msg.answer( await msg.answer(
"<b>Команды</b>\n\n" "<b>Команды</b>\n\n"
"<b>Сводки</b>\n" "<b>Сводки</b>\n"
"/servers — CPU, RAM, диск, аптайм по узлам\n" "/servers — CPU, RAM, диск, аптайм по узлам\n"
"/services — что работает и что нет\n" "/services — что работает и что нет\n"
"/certs — сколько дней осталось у TLS-сертификатов\n" "/certs — сколько дней осталось у TLS-сертификатов\n"
"/alerts — активные алерты сейчас\n\n" "/alerts — активные алерты сейчас\n\n"
"<b>Настройки</b>\n" "<b>Настройки</b>\n"
"/status — что включено\n" "/status — что включено\n"
"/alerts on | off — подписка на алерты" "/alerts on | off — подписка на алерты"
) )
@router.message(Command("status")) @router.message(Command("status"))
async def cmd_status(msg: Message): async def cmd_status(msg: Message):
row = db.get_user_info(msg.from_user.id) row = db.get_user_info(msg.from_user.id)
if not row: if not row:
await msg.answer("Нажми /start.") await msg.answer("Нажми /start.")
return return
await msg.answer("<b>Настройки</b>\n\n" + settings_block(row)) await msg.answer("<b>Настройки</b>\n\n" + settings_block(row))
@router.message(Command("alerts")) @router.message(Command("alerts"))
async def cmd_alerts(msg: Message, command: CommandObject): async def cmd_alerts(msg: Message, command: CommandObject):
uid = msg.from_user.id uid = msg.from_user.id
if not db.get_user_info(uid): if not db.get_user_info(uid):
db.add_user(uid, msg.from_user.username or str(uid)) db.add_user(uid, msg.from_user.username or str(uid))
arg = (command.args or "").strip().lower() arg = (command.args or "").strip().lower()
if arg in ("on", "вкл"): if arg in ("on", "вкл"):
db.set_alerts(uid, True) db.set_alerts(uid, True)
await msg.answer("Алерты включены.") await msg.answer("Алерты включены.")
return return
if arg in ("off", "выкл"): if arg in ("off", "выкл"):
db.set_alerts(uid, False) db.set_alerts(uid, False)
await msg.answer("Алерты выключены.") await msg.answer("Алерты выключены.")
return return
await msg.answer(await render_active_alerts()) await msg.answer(await render_active_alerts())
@router.message(Command("servers")) @router.message(Command("servers"))
async def cmd_servers(msg: Message): async def cmd_servers(msg: Message):
await msg.answer(await render_servers()) await msg.answer(await render_servers())
@router.message(Command("services")) @router.message(Command("services"))
async def cmd_services(msg: Message): async def cmd_services(msg: Message):
await msg.answer(await render_services()) await msg.answer(await render_services())
@router.message(Command("certs")) @router.message(Command("certs"))
async def cmd_certs(msg: Message): async def cmd_certs(msg: Message):
await msg.answer(await render_certs()) await msg.answer(await render_certs())
# ====================================================================== # ======================================================================
# Prometheus # Prometheus
# ====================================================================== # ======================================================================
async def _promq(session: aiohttp.ClientSession, query: str) -> List[dict]: def _auth(token: str) -> dict:
async with session.get(PROM_QUERY_URL, params={"query": query, "token": PROM_TOKEN}, """Token goes in the Authorization header, not in the URL (URLs end up in proxy logs)."""
timeout=aiohttp.ClientTimeout(total=ALERT_HTTP_TIMEOUT)) as r: return {"Authorization": f"Bearer {token}"}
r.raise_for_status()
j = await r.json()
if j.get("status") != "success": async def _promq(session: aiohttp.ClientSession, query: str) -> List[dict]:
raise RuntimeError(j.get("error", "prometheus error")) async with session.get(PROM_QUERY_URL, params={"query": query},
return j["data"]["result"] headers=_auth(PROM_TOKEN),
timeout=aiohttp.ClientTimeout(total=ALERT_HTTP_TIMEOUT)) as r:
r.raise_for_status()
def _by_server(result: List[dict]) -> Dict[str, float]: j = await r.json()
out = {} if j.get("status") != "success":
for s in result: raise RuntimeError(j.get("error", "prometheus error"))
srv = s["metric"].get("server") return j["data"]["result"]
if srv:
try:
out[srv] = float(s["value"][1]) def _by_server(result: List[dict]) -> Dict[str, float]:
except (ValueError, TypeError): out = {}
pass for s in result:
return out srv = s["metric"].get("server")
if srv:
try:
async def render_servers() -> str: out[srv] = float(s["value"][1])
q = { except (ValueError, TypeError):
"cpu": '100 - (avg by (server) (rate(node_cpu_seconds_total{mode="idle"}[2m])) * 100)', pass
"ram": '(1 - node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes) * 100', return out
"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', async def render_servers() -> str:
"up": 'up{job=~"node_.*"}', q = {
"boot": 'node_boot_time_seconds', "cpu": '100 - (avg by (server) (rate(node_cpu_seconds_total{mode="idle"}[2m])) * 100)',
"swap": '(node_memory_SwapTotal_bytes > bool 0) * (1 - node_memory_SwapFree_bytes / (node_memory_SwapTotal_bytes > 0)) * 100', "ram": '(1 - node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes) * 100',
# real hwmon sensors only (chip=~"pci.*") — excludes bogus acpitz/thermal_zone "disk": '(1 - avg by (server)(node_filesystem_avail_bytes{mountpoint="/",fstype!~"tmpfs|overlay|squashfs"}) '
# readings some laptops-as-servers report as a stuck fake value '/ avg by (server)(node_filesystem_size_bytes{mountpoint="/",fstype!~"tmpfs|overlay|squashfs"})) * 100',
"temp": 'max by (server) (node_hwmon_temp_celsius{chip=~"pci.*"})', "load": 'node_load1',
} "up": 'up{job=~"node_.*"}',
try: "boot": 'node_boot_time_seconds',
async with aiohttp.ClientSession() as s: "swap": '(node_memory_SwapTotal_bytes > bool 0) * (1 - node_memory_SwapFree_bytes / (node_memory_SwapTotal_bytes > 0)) * 100',
res = {k: _by_server(await _promq(s, v)) for k, v in q.items()} # real hwmon sensors only (chip=~"pci.*") — excludes bogus acpitz/thermal_zone
except Exception as e: # readings some laptops-as-servers report as a stuck fake value
return f"Не удалось получить метрики: {e}" "temp": 'max by (server) (node_hwmon_temp_celsius{chip=~"pci.*"})',
}
now = time.time() try:
out = ["<b>Серверы</b>"] async with aiohttp.ClientSession() as s:
for srv in SERVERS_ORDER: res = {k: _by_server(await _promq(s, v)) for k, v in q.items()}
if srv not in res["up"]: except Exception as e:
out.append(f"\n<b>{srv.upper()}</b> — нет данных") return f"Не удалось получить метрики: {e}"
continue
alive = res["up"].get(srv, 0) >= 1 now = time.time()
cpu, ram = res["cpu"].get(srv), res["ram"].get(srv) out = ["<b>Серверы</b>"]
disk, swap = res["disk"].get(srv), res["swap"].get(srv, 0.0) for srv in SERVERS_ORDER:
load = res["load"].get(srv) if srv not in res["up"]:
temp = res["temp"].get(srv) out.append(f"\n<b>{srv.upper()}</b> — нет данных")
up = fmt_uptime(now - res["boot"][srv]) if srv in res["boot"] else "?" continue
alive = res["up"].get(srv, 0) >= 1
def g(x): cpu, ram = res["cpu"].get(srv), res["ram"].get(srv)
return f"{x:.0f}%" if isinstance(x, (int, float)) else "?" disk, swap = res["disk"].get(srv), res["swap"].get(srv, 0.0)
load = res["load"].get(srv)
warn = [] temp = res["temp"].get(srv)
if isinstance(disk, float) and disk >= 85: up = fmt_uptime(now - res["boot"][srv]) if srv in res["boot"] else "?"
warn.append("диск")
if isinstance(ram, float) and ram >= 90: def g(x):
warn.append("память") return f"{x:.0f}%" if isinstance(x, (int, float)) else "?"
if isinstance(swap, float) and swap >= 60:
warn.append("swap") warn = []
if isinstance(temp, float) and temp >= 85: if isinstance(disk, float) and disk >= 85:
warn.append("температура") warn.append("диск")
head = f"\n<b>{srv.upper()}</b>" if isinstance(ram, float) and ram >= 90:
if not alive: warn.append("память")
head += " · НЕ ОТВЕЧАЕТ" if isinstance(swap, float) and swap >= 60:
elif warn: warn.append("swap")
head += " · ⚠ " + ", ".join(warn) if isinstance(temp, float) and temp >= 85:
out.append(head) warn.append("температура")
line = f"CPU {g(cpu)} · RAM {g(ram)} · диск {g(disk)} · swap {g(swap)}" head = f"\n<b>{srv.upper()}</b>"
if isinstance(temp, float): if not alive:
line += f" · temp {temp:.0f}°C" head += " · НЕ ОТВЕЧАЕТ"
out.append(line) elif warn:
out.append(f"load {load:.2f} · аптайм {up}" if isinstance(load, float) head += " · ⚠ " + ", ".join(warn)
else f"аптайм {up}") out.append(head)
return "\n".join(out) line = f"CPU {g(cpu)} · RAM {g(ram)} · диск {g(disk)} · swap {g(swap)}"
if isinstance(temp, float):
line += f" · temp {temp:.0f}°C"
async def render_services() -> str: out.append(line)
try: out.append(f"load {load:.2f} · аптайм {up}" if isinstance(load, float)
async with aiohttp.ClientSession() as s: else f"аптайм {up}")
res = await _promq(s, "up") return "\n".join(out)
except Exception as e:
return f"Не удалось получить статус сервисов: {e}"
total = len(res) async def render_services() -> str:
down = [(m["metric"].get("job", "?"), m["metric"].get("instance", "?")) try:
for m in res if m["value"][1] != "1"] async with aiohttp.ClientSession() as s:
if not down: res = await _promq(s, "up")
return f"<b>Сервисы</b>\n\nВсе {total} в норме." except Exception as e:
lines = [f"<b>Сервисы</b>\n\n{total - len(down)}/{total} в норме. Не отвечают:"] return f"Не удалось получить статус сервисов: {e}"
for job, inst in sorted(down): total = len(res)
lines.append(f"• {html.escape(job)} — {html.escape(inst)}") down = [(m["metric"].get("job", "?"), m["metric"].get("instance", "?"))
return "\n".join(lines) for m in res if m["value"][1] != "1"]
if not down:
return f"<b>Сервисы</b>\n\nВсе {total} в норме."
async def render_certs() -> str: lines = [f"<b>Сервисы</b>\n\n{total - len(down)}/{total} в норме. Не отвечают:"]
try: for job, inst in sorted(down):
async with aiohttp.ClientSession() as s: lines.append(f"• {html.escape(job)} — {html.escape(inst)}")
res = await _promq(s, "(probe_ssl_earliest_cert_expiry - time()) / 86400") return "\n".join(lines)
except Exception as e:
return f"Не удалось получить данные по сертификатам: {e}"
rows = [] async def render_certs() -> str:
for m in res: try:
try: async with aiohttp.ClientSession() as s:
days = float(m["value"][1]) res = await _promq(s, "(probe_ssl_earliest_cert_expiry - time()) / 86400")
except (ValueError, TypeError): except Exception as e:
continue return f"Не удалось получить данные по сертификатам: {e}"
host = m["metric"].get("instance", "?").replace("https://", "").split("/")[0] rows = []
rows.append((days, host)) for m in res:
if not rows: try:
return "Данных по сертификатам нет." days = float(m["value"][1])
rows.sort() except (ValueError, TypeError):
soon = [f"{h} — {d:.0f} дн" for d, h in rows if d < 14] continue
body = "\n".join(f"{d:>4.0f} {h}" for d, h in rows) host = m["metric"].get("instance", "?").replace("https://", "").split("/")[0]
head = "<b>TLS-сертификаты</b>" rows.append((days, host))
if soon: if not rows:
head += "\n⚠ скоро истекают: " + "; ".join(soon) return "Данных по сертификатам нет."
return head + f"\n<pre>{body}</pre>" 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 = "<b>TLS-сертификаты</b>"
# Alertmanager if soon:
# ====================================================================== head += "\n⚠ скоро истекают: " + "; ".join(soon)
def _sev(labels): return labels.get("severity", "").lower() return head + f"\n<pre>{body}</pre>"
def _where(labels): return labels.get("server") or labels.get("instance") or labels.get("job") or ""
# ======================================================================
def _fmt_firing(a: dict) -> str: # Alertmanager
lb, an = a.get("labels", {}), a.get("annotations", {}) # ======================================================================
name = html.escape(lb.get("alertname", "alert")) def _sev(labels): return labels.get("severity", "").lower()
sub = " · ".join(x for x in (_sev(lb), html.escape(_where(lb))) if x) def _where(labels): return labels.get("server") or labels.get("instance") or labels.get("job") or ""
summ = html.escape(an.get("summary") or an.get("description") or "")
text = f"<b>Проблема · {name}</b>"
if sub: def _fmt_firing(a: dict) -> str:
text += f"\n{sub}" lb, an = a.get("labels", {}), a.get("annotations", {})
if summ: name = html.escape(lb.get("alertname", "alert"))
text += f"\n{summ}" sub = " · ".join(x for x in (_sev(lb), html.escape(_where(lb))) if x)
return text summ = html.escape(an.get("summary") or an.get("description") or "")
text = f"<b>Проблема · {name}</b>"
if sub:
def _fmt_resolved(name: str) -> str: text += f"\n{sub}"
return f"<b>В норме · {html.escape(name)}</b>" if summ:
text += f"\n{summ}"
return text
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, def _fmt_resolved(name: str) -> str:
timeout=aiohttp.ClientTimeout(total=ALERT_HTTP_TIMEOUT)) as r: return f"<b>В норме · {html.escape(name)}</b>"
r.raise_for_status()
return await r.json()
async def _fetch_alerts(session: aiohttp.ClientSession) -> List[dict]:
params = {"active": "true", "silenced": "false", "inhibited": "false"}
async def _broadcast(text: str): async with session.get(AM_ALERTS_URL, params=params, headers=_auth(AM_TOKEN),
for uid in db.get_alert_users(): timeout=aiohttp.ClientTimeout(total=ALERT_HTTP_TIMEOUT)) as r:
try: r.raise_for_status()
await bot.send_message(uid, text) return await r.json()
except Exception as e:
log.error("alert to %s failed: %s", uid, e)
async def _broadcast(text: str):
for uid in db.get_alert_users():
async def render_active_alerts() -> str: try:
try: await bot.send_message(uid, text)
async with aiohttp.ClientSession() as s: except Exception as e:
data = await _fetch_alerts(s) log.error("alert to %s failed: %s", uid, e)
except Exception as e:
return f"Не удалось получить алерты: {e}"
firing = [a for a in data if a.get("status", {}).get("state") == "active"] async def render_active_alerts() -> str:
if not firing: try:
return "Активных алертов нет." async with aiohttp.ClientSession() as s:
return f"<b>Активные алерты: {len(firing)}</b>\n\n" + "\n\n".join(_fmt_firing(a) for a in firing) data = await _fetch_alerts(s)
except Exception as e:
return f"Не удалось получить алерты: {e}"
async def alert_poller(): firing = [a for a in data if a.get("status", {}).get("state") == "active"]
log.info("alert poller started (every %ss)", ALERT_POLL_SECONDS) if not firing:
async with aiohttp.ClientSession() as session: return "Активных алертов нет."
while True: return f"<b>Активные алерты: {len(firing)}</b>\n\n" + "\n\n".join(_fmt_firing(a) for a in firing)
try:
data = await _fetch_alerts(session)
current = {} def _fmt_watchdog_down(error: str) -> str:
for a in data: return ("<b>Мониторинг недоступен</b>\nAlertmanager не отвечает, о новых проблемах бот сейчас не узнает.\n"
if a.get("status", {}).get("state") != "active": f"Последняя ошибка: {html.escape(error)}")
continue
fp = a.get("fingerprint")
if not fp: def _fmt_watchdog_up() -> str:
continue return "<b>Мониторинг снова доступен</b>"
current[fp] = a
if db.get_alert_status(fp) != "firing":
await _broadcast(_fmt_firing(a)) async def process_alerts(data: List[dict]) -> None:
db.upsert_alert(fp, "firing", a.get("labels", {}).get("alertname", "alert")) """One poll: announce new firing alerts and alerts that have resolved since the last poll."""
for fp, name in db.list_firing_fingerprints(): current = {}
if fp not in current: for a in data:
await _broadcast(_fmt_resolved(name)) if a.get("status", {}).get("state") != "active":
db.upsert_alert(fp, "resolved", name) continue
db.purge_old_resolved() fp = a.get("fingerprint")
except asyncio.CancelledError: if not fp:
break continue
except Exception as e: current[fp] = a
log.warning("alert poll failed: %s", e) if db.get_alert_status(fp) != "firing":
await asyncio.sleep(ALERT_POLL_SECONDS) 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)

View file

@ -20,6 +20,9 @@ AM_TOKEN = os.environ["AM_TOKEN"]
PROM_TOKEN = os.environ["PROM_TOKEN"] PROM_TOKEN = os.environ["PROM_TOKEN"]
ALERT_POLL_SECONDS = int(os.environ.get("ALERT_POLL_SECONDS", "45")) ALERT_POLL_SECONDS = int(os.environ.get("ALERT_POLL_SECONDS", "45"))
ALERT_HTTP_TIMEOUT = int(os.environ.get("ALERT_HTTP_TIMEOUT", "15")) 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)) bot = Bot(token=BOT_TOKEN, default=DefaultBotProperties(parse_mode=ParseMode.HTML))
dp = Dispatcher(storage=MemoryStorage()) dp = Dispatcher(storage=MemoryStorage())

236
sql.py
View file

@ -1,118 +1,118 @@
import sqlite3 import sqlite3
from typing import List, Optional, Tuple from typing import List, Optional, Tuple
class DatabaseManager: class DatabaseManager:
def __init__(self, db_path: str = "notifications.db"): def __init__(self, db_path: str = "notifications.db"):
self.db_path = db_path self.db_path = db_path
self.init_database() self.init_database()
def _conn(self): def _conn(self):
return sqlite3.connect(self.db_path) return sqlite3.connect(self.db_path)
def init_database(self): def init_database(self):
with self._conn() as conn: with self._conn() as conn:
conn.execute(''' conn.execute('''
CREATE TABLE IF NOT EXISTS users ( CREATE TABLE IF NOT EXISTS users (
user_id INTEGER PRIMARY KEY, user_id INTEGER PRIMARY KEY,
username TEXT, username TEXT,
alerts_enabled BOOLEAN DEFAULT 1, alerts_enabled BOOLEAN DEFAULT 1,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
) )
''') ''')
conn.execute(''' conn.execute('''
CREATE TABLE IF NOT EXISTS alert_seen ( CREATE TABLE IF NOT EXISTS alert_seen (
fingerprint TEXT PRIMARY KEY, fingerprint TEXT PRIMARY KEY,
status TEXT, status TEXT,
name TEXT, name TEXT,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
) )
''') ''')
conn.commit() conn.commit()
# ---------- users ---------- # ---------- users ----------
def add_user(self, user_id: int, username: str) -> bool: def add_user(self, user_id: int, username: str) -> bool:
try: try:
with self._conn() as conn: with self._conn() as conn:
conn.execute( conn.execute(
"INSERT INTO users (user_id, username) VALUES (?, ?) " "INSERT INTO users (user_id, username) VALUES (?, ?) "
"ON CONFLICT(user_id) DO UPDATE SET username=excluded.username", "ON CONFLICT(user_id) DO UPDATE SET username=excluded.username",
(user_id, username), (user_id, username),
) )
conn.commit() conn.commit()
return True return True
except sqlite3.Error: except sqlite3.Error:
return False return False
def set_alerts(self, user_id: int, value: bool) -> bool: def set_alerts(self, user_id: int, value: bool) -> bool:
try: try:
with self._conn() as conn: with self._conn() as conn:
conn.execute("UPDATE users SET alerts_enabled = ? WHERE user_id = ?", conn.execute("UPDATE users SET alerts_enabled = ? WHERE user_id = ?",
(1 if value else 0, user_id)) (1 if value else 0, user_id))
conn.commit() conn.commit()
return True return True
except sqlite3.Error: except sqlite3.Error:
return False return False
def get_user_info(self, user_id: int) -> Optional[Tuple]: def get_user_info(self, user_id: int) -> Optional[Tuple]:
"""(user_id, username, alerts_enabled)""" """(user_id, username, alerts_enabled)"""
try: try:
with self._conn() as conn: with self._conn() as conn:
return conn.execute( return conn.execute(
"SELECT user_id, username, alerts_enabled FROM users WHERE user_id = ?", "SELECT user_id, username, alerts_enabled FROM users WHERE user_id = ?",
(user_id,) (user_id,)
).fetchone() ).fetchone()
except sqlite3.Error: except sqlite3.Error:
return None return None
def get_alert_users(self) -> List[int]: def get_alert_users(self) -> List[int]:
try: try:
with self._conn() as conn: with self._conn() as conn:
return [r[0] for r in conn.execute( return [r[0] for r in conn.execute(
"SELECT user_id FROM users WHERE alerts_enabled = 1" "SELECT user_id FROM users WHERE alerts_enabled = 1"
).fetchall()] ).fetchall()]
except sqlite3.Error: except sqlite3.Error:
return [] return []
# ---------- alert dedup ---------- # ---------- alert dedup ----------
def get_alert_status(self, fingerprint: str) -> Optional[str]: def get_alert_status(self, fingerprint: str) -> Optional[str]:
try: try:
with self._conn() as conn: with self._conn() as conn:
r = conn.execute("SELECT status FROM alert_seen WHERE fingerprint = ?", r = conn.execute("SELECT status FROM alert_seen WHERE fingerprint = ?",
(fingerprint,)).fetchone() (fingerprint,)).fetchone()
return r[0] if r else None return r[0] if r else None
except sqlite3.Error: except sqlite3.Error:
return None return None
def upsert_alert(self, fingerprint: str, status: str, name: str): def upsert_alert(self, fingerprint: str, status: str, name: str):
try: try:
with self._conn() as conn: with self._conn() as conn:
conn.execute( conn.execute(
"INSERT INTO alert_seen (fingerprint, status, name, updated_at) " "INSERT INTO alert_seen (fingerprint, status, name, updated_at) "
"VALUES (?, ?, ?, CURRENT_TIMESTAMP) " "VALUES (?, ?, ?, CURRENT_TIMESTAMP) "
"ON CONFLICT(fingerprint) DO UPDATE SET status=excluded.status, " "ON CONFLICT(fingerprint) DO UPDATE SET status=excluded.status, "
"name=excluded.name, updated_at=CURRENT_TIMESTAMP", "name=excluded.name, updated_at=CURRENT_TIMESTAMP",
(fingerprint, status, name), (fingerprint, status, name),
) )
conn.commit() conn.commit()
except sqlite3.Error: except sqlite3.Error:
pass pass
def list_firing_fingerprints(self) -> List[Tuple[str, str]]: def list_firing_fingerprints(self) -> List[Tuple[str, str]]:
try: try:
with self._conn() as conn: with self._conn() as conn:
return conn.execute( return conn.execute(
"SELECT fingerprint, name FROM alert_seen WHERE status = 'firing'" "SELECT fingerprint, name FROM alert_seen WHERE status = 'firing'"
).fetchall() ).fetchall()
except sqlite3.Error: except sqlite3.Error:
return [] return []
def purge_old_resolved(self, days: int = 3): def purge_old_resolved(self, days: int = 3):
try: try:
with self._conn() as conn: with self._conn() as conn:
conn.execute( conn.execute(
"DELETE FROM alert_seen WHERE status = 'resolved' " "DELETE FROM alert_seen WHERE status = 'resolved' "
"AND updated_at < datetime('now', ?)", (f'-{days} days',)) "AND updated_at < datetime('now', ?)", (f'-{days} days',))
conn.commit() conn.commit()
except sqlite3.Error: except sqlite3.Error:
pass pass