From a6968ebbe677bd1d72caf2cc741abc37c498ae03 Mon Sep 17 00:00:00 2001 From: h Date: Thu, 23 Jul 2026 03:45:08 +0200 Subject: [PATCH] feat(bot): show actual response, flap threshold, pause/rename/baseline, komodo fleet and alerts --- .env.example | 27 ++++ src/healthbot/__init__.py | 7 + src/healthbot/alerts.py | 106 +++++++++++++ src/healthbot/checker.py | 129 +++++++++++++--- src/healthbot/filters/admin.py | 4 +- src/healthbot/handlers/__init__.py | 4 +- src/healthbot/handlers/add.py | 8 +- src/healthbot/handlers/commands.py | 213 ++++++++++++++++++++++++--- src/healthbot/handlers/fleet.py | 205 ++++++++++++++++++++++++++ src/healthbot/keyboards/__init__.py | 33 ++++- src/healthbot/komodo.py | 132 +++++++++++++++++ src/healthbot/monitor.py | 120 +++++++++------ src/healthbot/notify.py | 28 ++++ src/healthbot/web.py | 69 +++++++++ src/utils/db/__init__.py | 17 +++ src/utils/db/models.py | 16 +- src/utils/db/repositories/service.py | 44 +++++- src/utils/env.py | 46 +++++- 18 files changed, 1105 insertions(+), 103 deletions(-) create mode 100644 src/healthbot/alerts.py create mode 100644 src/healthbot/handlers/fleet.py create mode 100644 src/healthbot/komodo.py create mode 100644 src/healthbot/notify.py create mode 100644 src/healthbot/web.py diff --git a/.env.example b/.env.example index 6827e93..a6afb21 100644 --- a/.env.example +++ b/.env.example @@ -1,12 +1,39 @@ COMPOSE_PROFILES=bot RUN_ENVIRONMENT=prod +# Кому писать. ADMIN_ID — один, ADMIN_IDS — список; работают вместе. ADMIN_ID= +# ADMIN_IDS=[111111111,222222222] + BOT__TOKEN= MONITOR__CHECK_INTERVAL=30 MONITOR__REMINDER_INTERVAL=1800 MONITOR__REQUEST_TIMEOUT=10 +# Сколько проверок подряд должны провалиться, прежде чем объявлять падение. +# 1 — будит по любому сетевому чиху, 2 — разумный минимум. +MONITOR__FAILURE_THRESHOLD=2 +# Сколько символов ответа показывать в уведомлении. +MONITOR__BODY_PREVIEW=400 + +## --- Komodo (необязательно) --------------------------------------------- +## Управление флотом из телеграма: /fleet, кнопки Deploy/Restart/Stop, логи. +## Ключ создаётся в Komodo: Settings -> API Keys. Права наследуются от +## пользователя, которым он создан, — заведи боту отдельного. +# KOMODO__URL=https://komodo.example.com +# KOMODO__KEY= +# KOMODO__SECRET= + +## --- Приём алертов из Komodo (необязательно) ---------------------------- +## Поднимает HTTP-приёмник; в Komodo заводится Alerter с endpoint type +## Custom и URL http://healthbot:8080/komodo/alert +## Порт лучше НЕ выставлять наружу: если Komodo и бот в одной docker-сети, +## публичный адрес не нужен вообще. +# WEB__ENABLED=true +# WEB__PORT=8080 +## Общий секрет. Komodo произвольные заголовки слать не умеет, поэтому его +## же можно дописать прямо в URL: .../komodo/alert?token=... +# WEB__TOKEN= DB__PATH=data/healthbot.db diff --git a/src/healthbot/__init__.py b/src/healthbot/__init__.py index f86a200..0e8e038 100644 --- a/src/healthbot/__init__.py +++ b/src/healthbot/__init__.py @@ -1,6 +1,7 @@ import asyncio import contextlib +from utils.env import env from utils.logging import logger, setup_logging setup_logging() @@ -11,19 +12,25 @@ async def runner() -> None: from . import handlers # noqa: PLC0415 from .common import bot, dp # noqa: PLC0415 + from .komodo import komodo # noqa: PLC0415 from .monitor import run_monitor # noqa: PLC0415 + from .web import run_web # noqa: PLC0415 await db.connect() dp.include_routers(handlers.router) await bot.delete_webhook(drop_pending_updates=True) monitor = asyncio.create_task(run_monitor()) + web = await run_web() if env.web.enabled else None try: await dp.start_polling(bot) finally: monitor.cancel() with contextlib.suppress(asyncio.CancelledError): await monitor + if web is not None: + await web.cleanup() + await komodo.close() await bot.session.close() await db.close() diff --git a/src/healthbot/alerts.py b/src/healthbot/alerts.py new file mode 100644 index 0000000..d0649c5 --- /dev/null +++ b/src/healthbot/alerts.py @@ -0,0 +1,106 @@ +"""Превращение алерта Komodo в сообщение для телеграма.""" + +import html +from typing import Any + +LEVEL_ICON = {"OK": "✅", "WARNING": "🟠", "CRITICAL": "🔴"} + + +def _gb(value: Any) -> str: # noqa: ANN401 + try: + return f"{float(value):.1f} ГБ" + except (TypeError, ValueError): + return str(value) + + +def _pct(value: Any) -> str: # noqa: ANN401 + try: + return f"{float(value):.1f}%" + except (TypeError, ValueError): + return str(value) + + +def _describe(kind: str, d: dict) -> tuple[str, list[str]]: # noqa: C901, PLR0911, PLR0912 + """Заголовок и строки подробностей для конкретного типа алерта.""" + name = d.get("name") or d.get("id") or "?" + server = d.get("server_name") + + if kind == "ServerUnreachable": + return f"Сервер {html.escape(name)} не отвечает", [ + f"Ошибка: {html.escape(str(d.get('err') or 'нет деталей'))}" + ] + if kind == "ServerCpu": + return f"CPU на {html.escape(name)}", [ + f"Загрузка: {_pct(d.get('percentage'))}" + ] + if kind == "ServerMem": + return f"Память на {html.escape(name)}", [ + f"Занято: {_gb(d.get('used_gb'))} из {_gb(d.get('total_gb'))}" + ] + if kind == "ServerDisk": + return f"Диск на {html.escape(name)}", [ + f"Точка: {html.escape(str(d.get('path', '/')))}", + f"Занято: {_gb(d.get('used_gb'))} из {_gb(d.get('total_gb'))}", + ] + if kind == "ServerVersionMismatch": + return f"Разъехались версии на {html.escape(name)}", [ + f"Periphery: {d.get('version')}, Core: {d.get('core_version')}", + "Чинится только по ssh - дашборд к такому серверу не достучится.", + ] + if kind in {"StackStateChange", "ContainerStateChange"}: + return f"{html.escape(name)}: {d.get('from')} → {d.get('to')}", ( + [f"Сервер: {html.escape(server)}"] if server else [] + ) + if kind in {"StackImageUpdateAvailable", "DeploymentImageUpdateAvailable"}: + return f"Обновление образа для {html.escape(name)}", [ + f"{html.escape(str(d.get('image', '?')))}" + ] + if kind in {"StackAutoUpdated", "DeploymentAutoUpdated"}: + return f"{html.escape(name)} обновлён автоматически", [ + f"{html.escape(str(d.get('images') or d.get('image') or ''))}" + ] + if kind == "ResourceSyncPendingUpdates": + return f"Синк {html.escape(name)}: git разошёлся с Komodo", [ + "Кто-то правил ресурсы мимо git, либо пуш не доехал." + ] + if kind in {"ProcedureFailed", "ActionFailed", "BuildFailed", "RepoBuildFailed"}: + what = {"ProcedureFailed": "Процедура", "ActionFailed": "Действие"}.get( + kind, "Сборка" + ) + return f"{what} {html.escape(name)} упала", [] + if kind == "ScheduleRun": + return f"По расписанию запущено: {html.escape(name)}", [] + if kind == "Test": + return "Проверка связи с Komodo", ["Если ты это видишь - канал работает."] + if kind == "Custom": + return html.escape(str(d.get("title") or "Сообщение из Komodo")), [ + html.escape(str(d.get("body") or "")) + ] + + rest = [ + f"{html.escape(str(k))}: {html.escape(str(v))}" + for k, v in d.items() + if k not in {"id", "server_id", "swarm_id"} + ] + return f"{html.escape(kind)}: {html.escape(name)}", rest[:8] + + +def format_alert(payload: dict) -> str: + level = str(payload.get("level", "")).upper() + resolved = bool(payload.get("resolved")) or level == "OK" + data = payload.get("data") or {} + kind = str(data.get("type") or "Unknown") + inner = data.get("data") or {} + + icon = "✅" if resolved else LEVEL_ICON.get(level, "🔔") + title, details = _describe(kind, inner if isinstance(inner, dict) else {}) + + if resolved: + return "\n".join( + [f"{icon} Отбой: {title}", f"Komodo · {html.escape(kind)}"] + ) + + lines = [f"{icon} {title}"] + lines.extend(x for x in details if x) + lines.append(f"Komodo · {html.escape(kind)}") + return "\n".join(lines) diff --git a/src/healthbot/checker.py b/src/healthbot/checker.py index d4342b3..9882385 100644 --- a/src/healthbot/checker.py +++ b/src/healthbot/checker.py @@ -1,39 +1,76 @@ """HTTP-пробинг сервисов и логика сравнения ответа с эталоном.""" +import difflib +import html import json -from dataclasses import dataclass +import time +from dataclasses import dataclass, field import aiohttp from utils.db.models import MatchType, Service +from utils.env import env from utils.format import short _MISSING = object() +DIFF_LINES = 12 +DIFF_HEADER_LINES = 2 + @dataclass(slots=True) class Probe: - reachable: bool # удалось ли вообще получить ответ + reachable: bool status: int | None body: str error: str | None + latency_ms: int | None = None + + +@dataclass(slots=True) +class Verdict: + """Результат сравнения.""" + + healthy: bool + reason: str = "" + expected: str | None = None + actual: str | None = None + details: list[str] = field(default_factory=list) + + def report(self) -> str: + """Человекочитаемое объяснение для сообщения в телеграм.""" + lines: list[str] = [] + if self.expected is not None or self.actual is not None: + lines.append(f"Ожидали: {html.escape(self.expected or '-')}") + lines.append(f"Сейчас: {html.escape(self.actual or '-')}") + lines.extend(self.details) + return "\n".join(lines) async def probe(session: aiohttp.ClientSession, url: str, timeout: int) -> Probe: # noqa: ASYNC109 + started = time.monotonic() try: async with session.get( url, timeout=aiohttp.ClientTimeout(total=timeout), allow_redirects=True ) as resp: body = await resp.text() - return Probe(reachable=True, status=resp.status, body=body, error=None) + return Probe( + reachable=True, + status=resp.status, + body=body, + error=None, + latency_ms=int((time.monotonic() - started) * 1000), + ) except TimeoutError: - return Probe(reachable=False, status=None, body="", error="таймаут") + return Probe( + reachable=False, status=None, body="", error=f"таймаут ({timeout}с)" + ) except aiohttp.ClientError as exc: return Probe(reachable=False, status=None, body="", error=f"нет связи: {exc}") def dig(data: object, path: str) -> object: - """Достаёт значение по dot-пути (индексы списков — числами).""" + """Достаёт значение по dot-пути (индексы списков - числами).""" node = data for segment in path.split("."): if isinstance(node, dict): @@ -49,32 +86,88 @@ def dig(data: object, path: str) -> object: return node -def evaluate(service: Service, result: Probe) -> tuple[bool, str]: # noqa: PLR0911 - """Возвращает (жив ли сервис, причина падения).""" +def snippet(text: str, limit: int | None = None) -> list[str]: + """Кусок ответа в
, готовый к вставке в сообщение."""
+    limit = limit or env.monitor.body_preview
+    body = text.strip()
+    if not body:
+        return ["(пустой ответ)"]
+    clipped = body[:limit]
+    tail = "\n…" if len(body) > limit else ""
+    return [f"
{html.escape(clipped)}{tail}
"] + + +def body_diff(expected: str, actual: str) -> list[str]: + """Чем именно тело отличается от эталона.""" + diff = list( + difflib.unified_diff( + expected.strip().splitlines(), actual.strip().splitlines(), lineterm="", n=0 + ) + )[DIFF_HEADER_LINES:] + if not diff: + return ["Отличие в невидимых символах.", *snippet(actual)] + shown = diff[:DIFF_LINES] + text = "\n".join(shown) + more = f"\n… ещё {len(diff) - len(shown)} строк" if len(diff) > len(shown) else "" + return [f"
{html.escape(text)}{more}
"] + + +def evaluate(service: Service, result: Probe) -> Verdict: # noqa: PLR0911 + """Жив ли сервис - и если нет, то что именно он сейчас отдаёт.""" if not result.reachable: - return False, result.error or "недоступен" + return Verdict(healthy=False, reason=result.error or "недоступен") if service.match_type is MatchType.STATUS: if result.status != service.ok_status: - return False, f"HTTP {result.status} (ждали {service.ok_status})" - return True, "" + return Verdict( + healthy=False, + reason=f"HTTP {result.status}", + expected=f"HTTP {service.ok_status}", + actual=f"HTTP {result.status}", + details=snippet(result.body), + ) + return Verdict(healthy=True) if service.match_type is MatchType.BODY: if result.body.strip() != service.ok_body.strip(): - return False, "тело ответа изменилось" - return True, "" + return Verdict( + healthy=False, + reason="тело ответа изменилось", + details=body_diff(service.ok_body, result.body), + ) + return Verdict(healthy=True) - # MatchType.JSON try: data = json.loads(result.body) except ValueError: - return False, "ответ перестал быть JSON" + return Verdict( + healthy=False, + reason="ответ перестал быть JSON", + expected=f"{service.json_path} = {service.json_value}", + actual=f"не JSON, HTTP {result.status}", + details=snippet(result.body), + ) + value = dig(data, service.json_path or "") if value is _MISSING: - return False, f"нет поля {service.json_path}" - if _scalar_str(value) != service.json_value: - return False, f"{service.json_path} = {short(_scalar_str(value))}" - return True, "" + return Verdict( + healthy=False, + reason=f"нет поля {service.json_path}", + expected=f"{service.json_path} = {service.json_value}", + actual="поля нет в ответе", + details=snippet(result.body), + ) + + current = _scalar_str(value) + if current != service.json_value: + return Verdict( + healthy=False, + reason=f"{service.json_path} = {short(current)}", + expected=service.json_value, + actual=current, + details=snippet(result.body), + ) + return Verdict(healthy=True) def _scalar_str(value: object) -> str: diff --git a/src/healthbot/filters/admin.py b/src/healthbot/filters/admin.py index 1618db7..1d5ee1e 100644 --- a/src/healthbot/filters/admin.py +++ b/src/healthbot/filters/admin.py @@ -5,8 +5,10 @@ from utils.env import env class Admin(BaseFilter): + """Пускаем только своих.""" + async def __call__(self, event: TelegramObject) -> bool: user = None if isinstance(event, (Message, CallbackQuery)): user = event.from_user - return user is not None and user.id == env.admin_id + return user is not None and user.id in env.admins diff --git a/src/healthbot/handlers/__init__.py b/src/healthbot/handlers/__init__.py index b2a6ec0..b4d3c2d 100644 --- a/src/healthbot/handlers/__init__.py +++ b/src/healthbot/handlers/__init__.py @@ -1,8 +1,8 @@ from aiogram import Router -from . import add, commands +from . import add, commands, fleet router = Router() -router.include_routers(commands.router, add.router) +router.include_routers(commands.router, add.router, fleet.router) __all__ = ["router"] diff --git a/src/healthbot/handlers/add.py b/src/healthbot/handlers/add.py index 5e3a0bf..624a35e 100644 --- a/src/healthbot/handlers/add.py +++ b/src/healthbot/handlers/add.py @@ -83,7 +83,7 @@ def _choose_view(data: dict) -> tuple[str, types.InlineKeyboardMarkup]: if len(kids) > MAX_CHILDREN: lines.append(f"…показаны первые {MAX_CHILDREN} полей") else: - lines.append("Ответ не JSON — можно следить за статусом или телом целиком.") + lines.append("Ответ не JSON - можно следить за статусом или телом целиком.") if path: kb.button(text="⬆️ Назад", callback_data="add:up") @@ -100,7 +100,7 @@ async def on_add_menu(callback: types.CallbackQuery, state: FSMContext) -> None: if isinstance(callback.message, types.Message): await callback.message.edit_text( "Пришли ссылку на сервис " - "(например https://example.com/health) — " + "(например https://example.com/health) - " "я запрошу её и предложу выбрать, за чем следить.", reply_markup=back_to_menu(), ) @@ -112,7 +112,7 @@ async def on_url(message: types.Message, state: FSMContext) -> None: url = (message.text or "").strip() if await db.services.by_url(url) is not None: - await message.answer("Этот сервис уже под наблюдением. /list — список.") + await message.answer("Этот сервис уже под наблюдением. /list - список.") return note = await message.answer("Запрашиваю…") @@ -123,7 +123,7 @@ async def on_url(message: types.Message, state: FSMContext) -> None: await state.clear() await note.edit_text( f"Не смог достучаться: {html.escape(result.error or '')}.\n" - "Эталон снимаю только с живого сервиса — попробуй, когда он поднимется.", + "Эталон снимаю только с живого сервиса - попробуй, когда он поднимется.", reply_markup=back_to_menu(), ) return diff --git a/src/healthbot/handlers/commands.py b/src/healthbot/handlers/commands.py index cb49c27..5b2eac3 100644 --- a/src/healthbot/handlers/commands.py +++ b/src/healthbot/handlers/commands.py @@ -1,12 +1,18 @@ -"""Меню, список сервисов, статистика, удаление.""" +"""Меню, список сервисов, карточка, статистика, ручная проверка, удаление.""" import html +import json as jsonlib +import aiohttp from aiogram import Bot, F, Router, types from aiogram.filters import Command, CommandStart +from aiogram.fsm.context import FSMContext +from aiogram.fsm.state import State, StatesGroup +from healthbot.checker import _scalar_str, dig, evaluate, probe, snippet from healthbot.filters import Admin from healthbot.keyboards import ( + after_check, back_to_menu, confirm_delete, main_menu, @@ -14,7 +20,8 @@ from healthbot.keyboards import ( services_list, ) from utils.db import db -from utils.db.models import Service +from utils.db.models import MatchType, Service +from utils.env import env from utils.format import human_duration, now, parse from utils.logging import logger @@ -22,21 +29,28 @@ router = Router() router.message.filter(Admin()) router.callback_query.filter(Admin()) +MAX_NAME = 64 + GREETING = ( - "👋 healthbot — слежу за твоими сервисами и пишу, если что-то легло.\n\n" + "👋 healthbot - слежу за твоими сервисами и пишу, если что-то легло.\n\n" "Пришли ссылку на health-эндпоинт или жми кнопку 👇" ) +class RenameFlow(StatesGroup): + waiting_name = State() + + @router.startup() async def on_startup(bot: Bot) -> None: - await bot.set_my_commands( - [ - types.BotCommand(command="start", description="Меню"), - types.BotCommand(command="list", description="Сервисы"), - types.BotCommand(command="stats", description="Статистика"), - ] - ) + commands = [ + types.BotCommand(command="start", description="Меню"), + types.BotCommand(command="list", description="Сервисы"), + types.BotCommand(command="stats", description="Статистика"), + ] + if env.komodo.enabled: + commands.append(types.BotCommand(command="fleet", description="Флот")) + await bot.set_my_commands(commands) logger.info(f"[green]Started as[/] @{(await bot.me()).username}") @@ -61,8 +75,11 @@ async def _list_view() -> tuple[str, types.InlineKeyboardMarkup]: services = await db.services.all() if not services: return "Пока ни одного сервиса. Пришли ссылку, чтобы добавить.", back_to_menu() - up = sum(s.is_up for s in services) - text = f"📋 Сервисы — 🟢 {up} / 🔴 {len(services) - up}\nТыкни, чтобы открыть." + watched = [s for s in services if not s.is_paused] + up = sum(s.is_up for s in watched) + paused = len(services) - len(watched) + tail = f", ⏸ {paused}" if paused else "" + text = f"📋 Сервисы - 🟢 {up} / 🔴 {len(watched) - up}{tail}\nТыкни, чтобы открыть." return text, services_list(services) @@ -94,8 +111,15 @@ async def _card_view(service: Service) -> str: window = (now() - parse(service.created_at)).total_seconds() uptime = 100.0 if window <= 0 else max(0.0, (window - down) / window * 100) - if service.is_up: + if service.is_paused: + status = "⏸ на паузе - не проверяется" + elif service.is_up: status = "🟢 работает" + if service.fail_streak: + status += ( + f"\n ⚠️ неудач подряд: {service.fail_streak}" + f"/{env.monitor.failure_threshold}" + ) else: current = await db.incidents.current(service.id) since = ( @@ -106,31 +130,176 @@ async def _card_view(service: Service) -> str: reason = html.escape(current.reason) if current else "" status = f"🔴 лежит уже {since}\n причина: {reason}" + latency = ( + f"{service.last_latency_ms} мс" if service.last_latency_ms is not None else "-" + ) + return ( f"{html.escape(service.name)}\n" f"{html.escape(service.url)}\n\n" f"Правило: {service.rule}\n" - f"Статус: {status}\n\n" + f"Статус: {status}\n" + f"Отклик: {latency}\n\n" f"📊 Аптайм: {uptime:.2f}%\n" f"Инцидентов: {len(incidents)}\n" f"Суммарный даунтайм: {human_duration(down) if down else '0с'}" ) -@router.callback_query(F.data.startswith("svc:open:")) -async def svc_open(callback: types.CallbackQuery) -> None: +async def _show_card(callback: types.CallbackQuery, service: Service) -> None: + if isinstance(callback.message, types.Message): + await callback.message.edit_text( + await _card_view(service), reply_markup=service_card(service) + ) + + +async def _load(callback: types.CallbackQuery) -> Service | None: service_id = int((callback.data or "").rsplit(":", 1)[-1]) service = await db.services.get(service_id) if service is None: await callback.answer("Сервис не найден", show_alert=True) + return service + + +@router.callback_query(F.data.startswith("svc:open:")) +async def svc_open(callback: types.CallbackQuery) -> None: + service = await _load(callback) + if service is None: return + await _show_card(callback, service) + await callback.answer() + + +@router.callback_query(F.data.startswith("svc:check:")) +async def svc_check(callback: types.CallbackQuery) -> None: + """Проверить прямо сейчас и показать, что сервис отдаёт.""" + service = await _load(callback) + if service is None: + return + await callback.answer("Запрашиваю…") + + async with aiohttp.ClientSession() as session: + result = await probe(session, service.url, env.monitor.request_timeout) + verdict = evaluate(service, result) + + head = "✅ Правило выполняется" if verdict.healthy else "🔴 Правило нарушено" + lines = [ + f"{html.escape(service.name)}", + f"{html.escape(service.url)}", + "", + head, + f"HTTP: {result.status if result.reachable else '-'}" + f" отклик: {result.latency_ms if result.latency_ms is not None else '-'} мс", + ] + if not verdict.healthy: + lines.append(f"Причина: {html.escape(verdict.reason)}") + report = verdict.report() + if report: + lines.append(report) + elif result.reachable: + lines.extend(snippet(result.body)) + if isinstance(callback.message, types.Message): await callback.message.edit_text( - await _card_view(service), reply_markup=service_card(service_id) + "\n".join(lines), + reply_markup=after_check( + service.id, offer_baseline=not verdict.healthy and result.reachable + ), + ) + + +@router.callback_query(F.data.startswith("svc:base:")) +async def svc_baseline(callback: types.CallbackQuery) -> None: + """Принять текущий ответ за новый эталон.""" + service = await _load(callback) + if service is None: + return + async with aiohttp.ClientSession() as session: + result = await probe(session, service.url, env.monitor.request_timeout) + if not result.reachable: + await callback.answer("Сервис недоступен - эталон снимать не с чего", True) # noqa: FBT003 + return + + json_value = service.json_value + if service.match_type is MatchType.JSON: + try: + value = dig(jsonlib.loads(result.body), service.json_path or "") + except ValueError: + await callback.answer("Ответ не JSON - эталон не обновлён", True) # noqa: FBT003 + return + json_value = _scalar_str(value) + + await db.services.set_baseline( + service.id, + ok_status=result.status or service.ok_status, + ok_body=( + result.body if service.match_type is MatchType.BODY else service.ok_body + ), + json_value=json_value, + ) + await db.services.set_streak(service.id, 0) + updated = await db.services.get(service.id) + if updated is not None: + await _show_card(callback, updated) + await callback.answer("Эталон обновлён") + + +@router.callback_query(F.data.startswith("svc:pause:")) +async def svc_pause(callback: types.CallbackQuery) -> None: + service = await _load(callback) + if service is None: + return + await db.services.set_paused(service.id, is_paused=True) + updated = await db.services.get(service.id) + if updated is not None: + await _show_card(callback, updated) + await callback.answer("На паузе") + + +@router.callback_query(F.data.startswith("svc:resume:")) +async def svc_resume(callback: types.CallbackQuery) -> None: + service = await _load(callback) + if service is None: + return + await db.services.set_paused(service.id, is_paused=False) + updated = await db.services.get(service.id) + if updated is not None: + await _show_card(callback, updated) + await callback.answer("Снова слежу") + + +@router.callback_query(F.data.startswith("svc:rename:")) +async def svc_rename(callback: types.CallbackQuery, state: FSMContext) -> None: + service = await _load(callback) + if service is None: + return + await state.set_state(RenameFlow.waiting_name) + await state.set_data({"service_id": service.id}) + if isinstance(callback.message, types.Message): + await callback.message.edit_text( + f"Как назвать вместо {html.escape(service.name)}?\n" + "Пришли новое имя сообщением.", + reply_markup=back_to_menu(), ) await callback.answer() +@router.message(RenameFlow.waiting_name) +async def on_new_name(message: types.Message, state: FSMContext) -> None: + name = (message.text or "").strip()[:MAX_NAME] + if not name: + await message.answer("Пустое имя не годится.") + return + data = await state.get_data() + await state.clear() + await db.services.rename(data["service_id"], name) + service = await db.services.get(data["service_id"]) + if service is None: + await message.answer("Сервис успел удалиться.") + return + await message.answer(await _card_view(service), reply_markup=service_card(service)) + + @router.callback_query(F.data.startswith("svc:del:")) async def svc_del(callback: types.CallbackQuery) -> None: service_id = int((callback.data or "").rsplit(":", 1)[-1]) @@ -162,11 +331,15 @@ async def _stats_view() -> tuple[str, types.InlineKeyboardMarkup]: down = _downtime(incidents) window = (now() - parse(service.created_at)).total_seconds() uptime = 100.0 if window <= 0 else max(0.0, (window - down) / window * 100) - mark = "🟢" if service.is_up else "🔴" + latency = ( + f", отклик {service.last_latency_ms} мс" + if service.last_latency_ms is not None + else "" + ) lines.append( - f"{mark} {html.escape(service.name)} — " + f"{service.mark} {html.escape(service.name)} - " f"аптайм {uptime:.2f}%, инцидентов {len(incidents)}, " - f"даунтайм {human_duration(down) if down else '0с'}" + f"даунтайм {human_duration(down) if down else '0с'}{latency}" ) return "\n".join(lines), back_to_menu() diff --git a/src/healthbot/handlers/fleet.py b/src/healthbot/handlers/fleet.py new file mode 100644 index 0000000..b3b9ac8 --- /dev/null +++ b/src/healthbot/handlers/fleet.py @@ -0,0 +1,205 @@ +"""Флот из телеграма: посмотреть сервера и стеки, нажать Deploy с телефона.""" + +import html + +from aiogram import F, Router, types +from aiogram.filters import Command +from aiogram.utils.keyboard import InlineKeyboardBuilder + +from healthbot.filters import Admin +from healthbot.keyboards import back_to_menu +from healthbot.komodo import KomodoError, Stack, komodo +from utils.logging import logger + +router = Router() +router.message.filter(Admin()) +router.callback_query.filter(Admin()) + +LOG_TAIL = 30 +LOG_CHARS = 2500 + +ACTIONS: dict[str, tuple[str, str, str]] = { + "dep": ("🚀 Deploy", "DeployStack", "Деплой запущен"), + "res": ("🔄 Restart", "RestartStack", "Перезапуск запущен"), + "sto": ("⏹ Stop", "StopStack", "Остановка запущена"), +} + + +async def _guard(callback: types.CallbackQuery) -> bool: + if komodo.enabled: + return True + await callback.answer("Komodo не настроен", show_alert=True) + return False + + +def _fleet_kb( + stacks: list[Stack], servers: dict[str, str] +) -> types.InlineKeyboardMarkup: + kb = InlineKeyboardBuilder() + for stack in sorted(stacks, key=lambda s: (servers.get(s.server_id, ""), s.name)): + kb.button( + text=f"{stack.mark} {stack.name}", callback_data=f"fleet:stack:{stack.id}" + ) + kb.button(text="🔄 Обновить", callback_data="fleet:main") + kb.button(text="⬅️ Меню", callback_data="menu:main") + kb.adjust(1) + return kb.as_markup() + + +async def _fleet_view() -> tuple[str, types.InlineKeyboardMarkup]: + servers = await komodo.servers() + stacks = await komodo.stacks() + by_id = {s.id: s.name for s in servers} + + lines = ["🛰 Флот", ""] + for server in sorted(servers, key=lambda s: s.name): + mine = [s for s in stacks if s.server_id == server.id] + good = sum(s.state == "running" for s in mine) + lines.append( + f"{server.mark} {html.escape(server.name)} - " + f"{good}/{len(mine)} стеков в строю" + ) + + broken = [s for s in stacks if s.state not in {"running", "unknown"}] + if broken: + lines.append("") + lines.append("Не в строю:") + lines.extend( + f"{s.mark} {html.escape(s.name)} - {html.escape(s.state)}" for s in broken + ) + + return "\n".join(lines), _fleet_kb(stacks, by_id) + + +async def _stack_view(stack_id: str) -> tuple[str, types.InlineKeyboardMarkup]: + stacks = await komodo.stacks() + stack = next((s for s in stacks if s.id == stack_id), None) + if stack is None: + return "Стек не найден - возможно, его удалили.", back_to_menu() + + servers = {s.id: s.name for s in await komodo.servers()} + text = ( + f"{stack.mark} {html.escape(stack.name)}\n" + f"Состояние: {html.escape(stack.state or '?')}\n" + f"Сервер: {html.escape(servers.get(stack.server_id, '?'))}" + ) + + kb = InlineKeyboardBuilder() + for key, (label, _, _) in ACTIONS.items(): + kb.button(text=label, callback_data=f"fleet:ask:{key}:{stack.id}") + kb.button(text="📜 Логи", callback_data=f"fleet:log:{stack.id}") + kb.button(text="⬅️ К флоту", callback_data="fleet:main") + kb.adjust(3, 2) + return text, kb.as_markup() + + +@router.message(Command("fleet")) +async def cmd_fleet(message: types.Message) -> None: + if not komodo.enabled: + await message.answer("Komodo не настроен - см. README, раздел про KOMODO__*.") + return + try: + text, markup = await _fleet_view() + except KomodoError as exc: + await message.answer(f"⚠️ {html.escape(str(exc))}") + return + await message.answer(text, reply_markup=markup) + + +@router.callback_query(F.data == "fleet:main") +async def fleet_main(callback: types.CallbackQuery) -> None: + if not await _guard(callback): + return + try: + text, markup = await _fleet_view() + except KomodoError as exc: + await callback.answer(str(exc), show_alert=True) + return + if isinstance(callback.message, types.Message): + await callback.message.edit_text(text, reply_markup=markup) + await callback.answer() + + +@router.callback_query(F.data.startswith("fleet:stack:")) +async def fleet_stack(callback: types.CallbackQuery) -> None: + if not await _guard(callback): + return + stack_id = (callback.data or "").rsplit(":", 1)[-1] + try: + text, markup = await _stack_view(stack_id) + except KomodoError as exc: + await callback.answer(str(exc), show_alert=True) + return + if isinstance(callback.message, types.Message): + await callback.message.edit_text(text, reply_markup=markup) + await callback.answer() + + +@router.callback_query(F.data.startswith("fleet:ask:")) +async def fleet_ask(callback: types.CallbackQuery) -> None: + """Подтверждение. Deploy с телефона слишком легко нажать штаниной.""" + _, _, key, stack_id = (callback.data or "").split(":", 3) + label = ACTIONS[key][0] + kb = InlineKeyboardBuilder() + kb.button(text=f"✅ Да, {label}", callback_data=f"fleet:do:{key}:{stack_id}") + kb.button(text="⬅️ Отмена", callback_data=f"fleet:stack:{stack_id}") + kb.adjust(1) + if isinstance(callback.message, types.Message): + await callback.message.edit_text( + f"{label} - подтверди, это подействует на прод.", + reply_markup=kb.as_markup(), + ) + await callback.answer() + + +@router.callback_query(F.data.startswith("fleet:do:")) +async def fleet_do(callback: types.CallbackQuery) -> None: + if not await _guard(callback): + return + _, _, key, stack_id = (callback.data or "").split(":", 3) + _, method, done = ACTIONS[key] + + await callback.answer(done) + if isinstance(callback.message, types.Message): + await callback.message.edit_text(f"⏳ {done}…") + + try: + await komodo.execute(method, {"stack": stack_id, "services": []}) + logger.info(f"{method} {stack_id} - запущено из телеграма") + except KomodoError as exc: + if isinstance(callback.message, types.Message): + await callback.message.edit_text( + f"⚠️ Не вышло: {html.escape(str(exc))}", reply_markup=back_to_menu() + ) + return + + try: + text, markup = await _stack_view(stack_id) + except KomodoError: + text, markup = "Готово.", back_to_menu() + if isinstance(callback.message, types.Message): + await callback.message.edit_text(f"✅ {done}.\n\n{text}", reply_markup=markup) + + +@router.callback_query(F.data.startswith("fleet:log:")) +async def fleet_log(callback: types.CallbackQuery) -> None: + if not await _guard(callback): + return + stack_id = (callback.data or "").rsplit(":", 1)[-1] + await callback.answer("Забираю логи…") + try: + log = await komodo.stack_log(stack_id, tail=LOG_TAIL) + except KomodoError as exc: + await callback.answer(str(exc), show_alert=True) + return + + kb = InlineKeyboardBuilder() + kb.button(text="🔄 Ещё раз", callback_data=f"fleet:log:{stack_id}") + kb.button(text="⬅️ К стеку", callback_data=f"fleet:stack:{stack_id}") + kb.adjust(2) + tail = log[-LOG_CHARS:] + if isinstance(callback.message, types.Message): + await callback.message.edit_text( + f"📜 Последние {LOG_TAIL} строк:\n
{html.escape(tail)}
", + reply_markup=kb.as_markup(), + ) diff --git a/src/healthbot/keyboards/__init__.py b/src/healthbot/keyboards/__init__.py index 65c5e8b..867252b 100644 --- a/src/healthbot/keyboards/__init__.py +++ b/src/healthbot/keyboards/__init__.py @@ -2,6 +2,7 @@ from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup from aiogram.utils.keyboard import InlineKeyboardBuilder from utils.db.models import Service +from utils.env import env def main_menu() -> InlineKeyboardMarkup: @@ -9,26 +10,35 @@ def main_menu() -> InlineKeyboardMarkup: kb.button(text="➕ Добавить сервис", callback_data="menu:add") kb.button(text="📋 Сервисы", callback_data="menu:list") kb.button(text="📊 Статистика", callback_data="menu:stats") - kb.adjust(1, 2) + if env.komodo.enabled: + kb.button(text="🛰 Флот", callback_data="fleet:main") + kb.adjust(1, 2, 1) + else: + kb.adjust(1, 2) return kb.as_markup() def services_list(services: list[Service]) -> InlineKeyboardMarkup: kb = InlineKeyboardBuilder() for svc in services: - mark = "🟢" if svc.is_up else "🔴" - kb.button(text=f"{mark} {svc.name}", callback_data=f"svc:open:{svc.id}") + kb.button(text=f"{svc.mark} {svc.name}", callback_data=f"svc:open:{svc.id}") kb.button(text="➕ Добавить", callback_data="menu:add") kb.button(text="⬅️ Меню", callback_data="menu:main") kb.adjust(1) return kb.as_markup() -def service_card(service_id: int) -> InlineKeyboardMarkup: +def service_card(service: Service) -> InlineKeyboardMarkup: kb = InlineKeyboardBuilder() - kb.button(text="🗑 Удалить", callback_data=f"svc:del:{service_id}") + kb.button(text="🔍 Проверить сейчас", callback_data=f"svc:check:{service.id}") + if service.is_paused: + kb.button(text="▶️ Возобновить", callback_data=f"svc:resume:{service.id}") + else: + kb.button(text="⏸ Пауза", callback_data=f"svc:pause:{service.id}") + kb.button(text="✏️ Переименовать", callback_data=f"svc:rename:{service.id}") + kb.button(text="🗑 Удалить", callback_data=f"svc:del:{service.id}") kb.button(text="⬅️ К списку", callback_data="menu:list") - kb.adjust(2) + kb.adjust(1, 2, 2) return kb.as_markup() @@ -40,6 +50,17 @@ def confirm_delete(service_id: int) -> InlineKeyboardMarkup: return kb.as_markup() +def after_check(service_id: int, *, offer_baseline: bool) -> InlineKeyboardMarkup: + """Клавиатура под результатом ручной проверки.""" + kb = InlineKeyboardBuilder() + if offer_baseline: + kb.button(text="📌 Принять как эталон", callback_data=f"svc:base:{service_id}") + kb.button(text="🔄 Ещё раз", callback_data=f"svc:check:{service_id}") + kb.button(text="⬅️ К сервису", callback_data=f"svc:open:{service_id}") + kb.adjust(1, 2) + return kb.as_markup() + + def back_to_menu() -> InlineKeyboardMarkup: return InlineKeyboardMarkup( inline_keyboard=[ diff --git a/src/healthbot/komodo.py b/src/healthbot/komodo.py new file mode 100644 index 0000000..e8a970f --- /dev/null +++ b/src/healthbot/komodo.py @@ -0,0 +1,132 @@ +"""Тонкий клиент Komodo: посмотреть флот и нажать кнопку.""" + +import re +from dataclasses import dataclass +from typing import Any + +import aiohttp + +from utils.env import env +from utils.logging import logger + +ANSI_RE = re.compile(r"\x1b\[[0-9;]*[a-zA-Z]") + +OK = 200 + + +def strip_ansi(text: str) -> str: + return ANSI_RE.sub("", text) + + +class KomodoError(RuntimeError): + pass + + +@dataclass(slots=True) +class Server: + id: str + name: str + state: str + + @property + def mark(self) -> str: + return {"Ok": "🟢", "NotOk": "🔴", "Disabled": "⏸"}.get(self.state, "⚪️") + + +@dataclass(slots=True) +class Stack: + id: str + name: str + state: str + server_id: str + + @property + def mark(self) -> str: + return { + "running": "🟢", + "unhealthy": "🟠", + "paused": "⏸", + "down": "⚫️", + "stopped": "🔴", + "restarting": "🔄", + "unknown": "⚪️", + }.get(self.state, "⚪️") + + +class Komodo: + def __init__(self) -> None: + self._session: aiohttp.ClientSession | None = None + + @property + def enabled(self) -> bool: + return env.komodo.enabled + + async def _call(self, route: str, kind: str, params: dict | None = None) -> Any: # noqa: ANN401 + if not self.enabled: + msg = "Komodo не настроен (KOMODO__URL / KOMODO__KEY / KOMODO__SECRET)" + raise KomodoError(msg) + if self._session is None or self._session.closed: + self._session = aiohttp.ClientSession() + url = env.komodo.url.rstrip("/") + route + try: + async with self._session.post( + url, + json={"type": kind, "params": params or {}}, + headers={ + "X-Api-Key": env.komodo.key.get_secret_value(), + "X-Api-Secret": env.komodo.secret.get_secret_value(), + }, + timeout=aiohttp.ClientTimeout(total=env.komodo.timeout), + ) as resp: + text = await resp.text() + if resp.status != OK: + logger.warning(f"Komodo {kind} -> HTTP {resp.status}: {text[:300]}") + msg = f"Komodo ответил HTTP {resp.status}" + raise KomodoError(msg) + return await resp.json() + except TimeoutError as exc: + msg = f"Komodo не ответил за {env.komodo.timeout}с" + raise KomodoError(msg) from exc + except aiohttp.ClientError as exc: + msg = f"Нет связи с Komodo: {exc}" + raise KomodoError(msg) from exc + + async def read(self, kind: str, params: dict | None = None) -> Any: # noqa: ANN401 + return await self._call("/read", kind, params) + + async def execute(self, kind: str, params: dict | None = None) -> Any: # noqa: ANN401 + return await self._call("/execute", kind, params) + + async def servers(self) -> list[Server]: + raw = await self.read("ListServers") + return [ + Server(id=s["id"], name=s["name"], state=s.get("info", {}).get("state", "")) + for s in raw + ] + + async def stacks(self) -> list[Stack]: + raw = await self.read("ListStacks") + return [ + Stack( + id=s["id"], + name=s["name"], + state=s.get("info", {}).get("state", ""), + server_id=s.get("info", {}).get("server_id", ""), + ) + for s in raw + ] + + async def stack_log(self, stack: str, tail: int = 40) -> str: + """Хвост лога стека одной строкой.""" + raw = await self.read( + "GetStackLog", {"stack": stack, "services": [], "tail": tail} + ) + out = (raw.get("stdout") or "") + (raw.get("stderr") or "") + return strip_ansi(out).strip() or "(пусто)" + + async def close(self) -> None: + if self._session is not None and not self._session.closed: + await self._session.close() + + +komodo = Komodo() diff --git a/src/healthbot/monitor.py b/src/healthbot/monitor.py index 573017d..91ab5bd 100644 --- a/src/healthbot/monitor.py +++ b/src/healthbot/monitor.py @@ -5,8 +5,8 @@ import html import aiohttp -from healthbot.checker import evaluate, probe -from healthbot.common import bot +from healthbot.checker import Verdict, evaluate, probe +from healthbot.notify import notify from utils.db import db from utils.db.models import Service from utils.env import env @@ -14,62 +14,96 @@ from utils.format import human_duration, now, parse from utils.logging import logger -async def _notify(text: str) -> None: - try: - await bot.send_message(env.admin_id, text) - except Exception: - logger.exception("Не смог отправить уведомление админу") +def _head(service: Service, icon: str, tail: str) -> str: + return ( + f"{icon} {html.escape(service.name)} {tail}\n" + f"{html.escape(service.url)}" + ) -async def _check_one(session: aiohttp.ClientSession, service: Service) -> None: +async def _on_recovered(service: Service) -> None: + current = await db.incidents.current(service.id) + downtime = "-" + if current is not None: + await db.incidents.close(current.id) + downtime = human_duration((now() - parse(current.started_at)).total_seconds()) + await db.services.set_up(service.id, is_up=True) + await notify(f"{_head(service, '✅', 'снова в строю')}\nДаунтайм: {downtime}") + logger.info(f"[green]UP[/] {service.name}") + + +async def _on_failed(service: Service, verdict: Verdict) -> None: + await db.incidents.open(service.id, verdict.reason) + await db.services.set_up(service.id, is_up=False) + lines = [_head(service, "🔴", "лёг"), f"Причина: {html.escape(verdict.reason)}"] + report = verdict.report() + if report: + lines.append(report) + await notify("\n".join(lines)) + logger.warning(f"[red]DOWN[/] {service.name}: {verdict.reason}") + + +async def _on_still_down(service: Service, verdict: Verdict) -> None: + current = await db.incidents.current(service.id) + if current is None: + return + idle = (now() - parse(current.last_reminder_at)).total_seconds() + if idle < env.monitor.reminder_interval: + return + await db.incidents.touch_reminder(current.id) + since = human_duration((now() - parse(current.started_at)).total_seconds()) + lines = [ + _head(service, "⏰", f"всё ещё лежит (уже {since})"), + f"Причина: {html.escape(verdict.reason)}", + ] + report = verdict.report() + if report: + lines.append(report) + await notify("\n".join(lines)) + + +async def check_one(session: aiohttp.ClientSession, service: Service) -> Verdict: + """Одна проверка со всеми побочными эффектами. Возвращает вердикт.""" result = await probe(session, service.url, env.monitor.request_timeout) - healthy, reason = evaluate(service, result) - name = html.escape(service.name) + verdict = evaluate(service, result) + await db.services.touch_probe(service.id, result.latency_ms) - if healthy and not service.is_up: - current = await db.incidents.current(service.id) - downtime = "" - if current is not None: - await db.incidents.close(current.id) - downtime = human_duration( - (now() - parse(current.started_at)).total_seconds() - ) - await db.services.set_up(service.id, is_up=True) - await _notify(f"✅ {name} снова в строю.\nДаунтайм: {downtime}.") - logger.info(f"[green]UP[/] {service.name}") + if verdict.healthy: + if service.fail_streak: + await db.services.set_streak(service.id, 0) + if not service.is_up: + await _on_recovered(service) + return verdict - elif not healthy and service.is_up: - await db.incidents.open(service.id, reason) - await db.services.set_up(service.id, is_up=False) - await _notify( - f"🔴 {name} лёг.\n" - f"{html.escape(service.url)}\n" - f"Причина: {html.escape(reason)}" - ) - logger.warning(f"[red]DOWN[/] {service.name}: {reason}") + streak = service.fail_streak + 1 + await db.services.set_streak(service.id, streak) - elif not healthy and not service.is_up: - current = await db.incidents.current(service.id) - if current is None: - return - idle = (now() - parse(current.last_reminder_at)).total_seconds() - if idle >= env.monitor.reminder_interval: - await db.incidents.touch_reminder(current.id) - since = human_duration((now() - parse(current.started_at)).total_seconds()) - await _notify( - f"⏰ {name} всё ещё лежит (уже {since}).\n" - f"Причина: {html.escape(reason)}" + if service.is_up: + if streak < env.monitor.failure_threshold: + logger.info( + f"[yellow]FLAP[/] {service.name}: " + f"{streak}/{env.monitor.failure_threshold} - {verdict.reason}" ) + return verdict + await _on_failed(service, verdict) + else: + await _on_still_down(service, verdict) + return verdict async def run_monitor() -> None: - logger.info(f"Monitor: каждые {env.monitor.check_interval}с") + logger.info( + f"Monitor: каждые {env.monitor.check_interval}с, " + f"порог падения {env.monitor.failure_threshold}" + ) async with aiohttp.ClientSession() as session: while True: try: for service in await db.services.all(): + if service.is_paused: + continue try: - await _check_one(session, service) + await check_one(session, service) except Exception: logger.exception(f"Ошибка проверки {service.url}") except Exception: diff --git a/src/healthbot/notify.py b/src/healthbot/notify.py new file mode 100644 index 0000000..a57b046 --- /dev/null +++ b/src/healthbot/notify.py @@ -0,0 +1,28 @@ +"""Отправка уведомлений админам.""" + +from healthbot.common import bot +from utils.env import env +from utils.logging import logger + +TELEGRAM_LIMIT = 4096 +SAFE_LIMIT = 3900 + + +def clip(text: str) -> str: + if len(text) <= TELEGRAM_LIMIT: + return text + cut = text[:SAFE_LIMIT] + cut = cut[: cut.rfind("\n")] if "\n" in cut else cut + if cut.count("
") > cut.count("
"): + cut += "
" + return cut + "\n…обрезано" + + +async def notify(text: str) -> None: + """Разослать всем админам. Падение одного не мешает остальным.""" + payload = clip(text) + for admin in env.admins: + try: + await bot.send_message(admin, payload) + except Exception: + logger.exception(f"Не смог отправить уведомление {admin}") diff --git a/src/healthbot/web.py b/src/healthbot/web.py new file mode 100644 index 0000000..dcbbde2 --- /dev/null +++ b/src/healthbot/web.py @@ -0,0 +1,69 @@ +"""HTTP-приёмник входящих алертов.""" + +import hmac +import json + +from aiohttp import web + +from healthbot.alerts import format_alert +from healthbot.notify import notify +from utils.env import env +from utils.logging import logger + +MUTED = frozenset({"ScheduleRun"}) + + +def _authorized(request: web.Request) -> bool: + expected = env.web.token.get_secret_value() + if not expected: + return True + given = request.headers.get("X-Alert-Token") or request.query.get("token") or "" + return hmac.compare_digest(given, expected) + + +async def handle_alert(request: web.Request) -> web.Response: + if not _authorized(request): + logger.warning(f"Отклонён алерт без токена от {request.remote}") + return web.json_response({"error": "unauthorized"}, status=401) + + raw = await request.text() + try: + payload = json.loads(raw) + except ValueError: + logger.warning(f"Алерт не разобрался как JSON: {raw[:200]}") + return web.json_response({"error": "bad json"}, status=400) + + if not isinstance(payload, dict): + return web.json_response({"error": "expected object"}, status=400) + + try: + text = format_alert(payload) + except Exception: + logger.exception("Не смог отформатировать алерт") + text = f"🔔 Алерт из Komodo (не разобрал формат)\n
{raw[:1000]}
" + + kind = (payload.get("data") or {}).get("type", "?") + logger.info(f"Алерт принят: {kind} (level={payload.get('level')})") + + if kind in MUTED: + return web.json_response({"ok": True, "muted": True}) + + await notify(text) + return web.json_response({"ok": True}) + + +async def handle_health(_: web.Request) -> web.Response: + return web.json_response({"status": "ok"}) + + +async def run_web() -> web.AppRunner: + app = web.Application() + app.router.add_post("/komodo/alert", handle_alert) + app.router.add_get("/health", handle_health) + + runner = web.AppRunner(app) + await runner.setup() + site = web.TCPSite(runner, env.web.host, env.web.port) + await site.start() + logger.info(f"Web: слушаю {env.web.host}:{env.web.port}") + return runner diff --git a/src/utils/db/__init__.py b/src/utils/db/__init__.py index 9e155d5..9342ffe 100644 --- a/src/utils/db/__init__.py +++ b/src/utils/db/__init__.py @@ -32,6 +32,13 @@ CREATE TABLE IF NOT EXISTS incidents ( CREATE INDEX IF NOT EXISTS idx_incidents_service ON incidents(service_id); """ +MIGRATIONS: tuple[tuple[str, str, str], ...] = ( + ("services", "fail_streak", "INTEGER NOT NULL DEFAULT 0"), + ("services", "is_paused", "INTEGER NOT NULL DEFAULT 0"), + ("services", "last_latency_ms", "INTEGER"), + ("services", "last_checked_at", "TEXT"), +) + class Database: def __init__(self, path: str) -> None: @@ -46,10 +53,20 @@ class Database: self._conn.row_factory = aiosqlite.Row await self._conn.execute("PRAGMA foreign_keys = ON") await self._conn.executescript(SCHEMA) + await self._migrate() await self._conn.commit() self.services = ServiceRepository(self._conn) self.incidents = IncidentRepository(self._conn) + async def _migrate(self) -> None: + assert self._conn is not None + for table, column, decl in MIGRATIONS: + cursor = await self._conn.execute(f"PRAGMA table_info({table})") + existing = {row["name"] for row in await cursor.fetchall()} + if column in existing: + continue + await self._conn.execute(f"ALTER TABLE {table} ADD COLUMN {column} {decl}") + async def close(self) -> None: if self._conn is not None: await self._conn.close() diff --git a/src/utils/db/models.py b/src/utils/db/models.py index abc4917..782f630 100644 --- a/src/utils/db/models.py +++ b/src/utils/db/models.py @@ -3,9 +3,9 @@ from enum import StrEnum class MatchType(StrEnum): - STATUS = "status" # следим только за HTTP-статусом - BODY = "body" # тело ответа должно совпадать целиком - JSON = "json" # конкретное поле JSON должно быть равно эталону + STATUS = "status" + BODY = "body" + JSON = "json" @dataclass(slots=True) @@ -20,6 +20,10 @@ class Service: json_value: str | None is_up: bool created_at: str + fail_streak: int = 0 + is_paused: bool = False + last_latency_ms: int | None = None + last_checked_at: str | None = None @property def rule(self) -> str: @@ -29,6 +33,12 @@ class Service: return f"{self.json_path} = {self.json_value}" return "тело ответа неизменно" + @property + def mark(self) -> str: + if self.is_paused: + return "⏸" + return "🟢" if self.is_up else "🔴" + @dataclass(slots=True) class Incident: diff --git a/src/utils/db/repositories/service.py b/src/utils/db/repositories/service.py index 4106592..80db74a 100644 --- a/src/utils/db/repositories/service.py +++ b/src/utils/db/repositories/service.py @@ -5,7 +5,8 @@ from utils.format import now _COLUMNS = ( "id, name, url, match_type, ok_status, ok_body, " - "json_path, json_value, is_up, created_at" + "json_path, json_value, is_up, created_at, " + "fail_streak, is_paused, last_latency_ms, last_checked_at" ) @@ -21,6 +22,10 @@ def _row_to_service(row: aiosqlite.Row) -> Service: json_value=row["json_value"], is_up=bool(row["is_up"]), created_at=row["created_at"], + fail_streak=row["fail_streak"], + is_paused=bool(row["is_paused"]), + last_latency_ms=row["last_latency_ms"], + last_checked_at=row["last_checked_at"], ) @@ -87,6 +92,43 @@ class ServiceRepository: ) await self._conn.commit() + async def set_streak(self, service_id: int, streak: int) -> None: + await self._conn.execute( + "UPDATE services SET fail_streak = ? WHERE id = ?", (streak, service_id) + ) + await self._conn.commit() + + async def set_paused(self, service_id: int, *, is_paused: bool) -> None: + await self._conn.execute( + "UPDATE services SET is_paused = ?, fail_streak = 0 WHERE id = ?", + (int(is_paused), service_id), + ) + await self._conn.commit() + + async def touch_probe(self, service_id: int, latency_ms: int | None) -> None: + await self._conn.execute( + "UPDATE services SET last_latency_ms = ?, last_checked_at = ? WHERE id = ?", + (latency_ms, now().isoformat(), service_id), + ) + await self._conn.commit() + + async def rename(self, service_id: int, name: str) -> None: + await self._conn.execute( + "UPDATE services SET name = ? WHERE id = ?", (name, service_id) + ) + await self._conn.commit() + + async def set_baseline( + self, service_id: int, *, ok_status: int, ok_body: str, json_value: str | None + ) -> None: + """Переснять эталон с текущего ответа.""" + await self._conn.execute( + "UPDATE services SET ok_status = ?, ok_body = ?, json_value = ? " + "WHERE id = ?", + (ok_status, ok_body, json_value, service_id), + ) + await self._conn.commit() + async def delete(self, service_id: int) -> None: await self._conn.execute("DELETE FROM services WHERE id = ?", (service_id,)) await self._conn.commit() diff --git a/src/utils/env.py b/src/utils/env.py index 33b8ab6..68b4a20 100644 --- a/src/utils/env.py +++ b/src/utils/env.py @@ -1,4 +1,4 @@ -from pydantic import Field, SecretStr +from pydantic import Field, SecretStr, computed_field from pydantic_settings import BaseSettings, SettingsConfigDict @@ -7,9 +7,36 @@ class BotSettings(BaseSettings): class MonitorSettings(BaseSettings): - check_interval: int = 30 # как часто пинговать сервисы, сек - reminder_interval: int = 1800 # как часто напоминать что всё ещё лежит, сек - request_timeout: int = 10 # таймаут одного запроса, сек + check_interval: int = 30 + reminder_interval: int = 1800 + request_timeout: int = 10 + + failure_threshold: int = 2 + + body_preview: int = 400 + + +class KomodoSettings(BaseSettings): + """Связь с Komodo. Всё опционально - без неё бот работает как раньше.""" + + url: str = "" + key: SecretStr = SecretStr("") + secret: SecretStr = SecretStr("") + timeout: int = 20 + + @computed_field + @property + def enabled(self) -> bool: + return bool(self.url and self.key.get_secret_value()) + + +class WebSettings(BaseSettings): + """Приёмник входящих алертов (Komodo Alerter с endpoint type Custom).""" + + enabled: bool = False + host: str = "0.0.0.0" # noqa: S104 + port: int = 8080 + token: SecretStr = SecretStr("") class DbSettings(BaseSettings): @@ -24,9 +51,13 @@ class LogSettings(BaseSettings): class Settings(BaseSettings): - admin_id: int + admin_id: int = 0 + admin_ids: list[int] = Field(default_factory=list) + bot: BotSettings = Field(default_factory=BotSettings) monitor: MonitorSettings = Field(default_factory=MonitorSettings) + komodo: KomodoSettings = Field(default_factory=KomodoSettings) + web: WebSettings = Field(default_factory=WebSettings) db: DbSettings = Field(default_factory=DbSettings) log: LogSettings = Field(default_factory=LogSettings) @@ -34,5 +65,10 @@ class Settings(BaseSettings): case_sensitive=False, env_file=".env", env_nested_delimiter="__", extra="ignore" ) + @computed_field + @property + def admins(self) -> set[int]: + return {i for i in (self.admin_id, *self.admin_ids) if i} + env = Settings()