feat(bot): show actual response, flap threshold, pause/rename/baseline, komodo fleet and alerts

This commit is contained in:
hh
2026-07-23 03:45:08 +02:00
parent a46bd92b4c
commit a6968ebbe6
18 changed files with 1105 additions and 103 deletions
+27
View File
@@ -1,12 +1,39 @@
COMPOSE_PROFILES=bot COMPOSE_PROFILES=bot
RUN_ENVIRONMENT=prod RUN_ENVIRONMENT=prod
# Кому писать. ADMIN_ID — один, ADMIN_IDS — список; работают вместе.
ADMIN_ID=<ADMIN_ID> ADMIN_ID=<ADMIN_ID>
# ADMIN_IDS=[111111111,222222222]
BOT__TOKEN=<BOT__TOKEN> BOT__TOKEN=<BOT__TOKEN>
MONITOR__CHECK_INTERVAL=30 MONITOR__CHECK_INTERVAL=30
MONITOR__REMINDER_INTERVAL=1800 MONITOR__REMINDER_INTERVAL=1800
MONITOR__REQUEST_TIMEOUT=10 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 DB__PATH=data/healthbot.db
+7
View File
@@ -1,6 +1,7 @@
import asyncio import asyncio
import contextlib import contextlib
from utils.env import env
from utils.logging import logger, setup_logging from utils.logging import logger, setup_logging
setup_logging() setup_logging()
@@ -11,19 +12,25 @@ async def runner() -> None:
from . import handlers # noqa: PLC0415 from . import handlers # noqa: PLC0415
from .common import bot, dp # noqa: PLC0415 from .common import bot, dp # noqa: PLC0415
from .komodo import komodo # noqa: PLC0415
from .monitor import run_monitor # noqa: PLC0415 from .monitor import run_monitor # noqa: PLC0415
from .web import run_web # noqa: PLC0415
await db.connect() await db.connect()
dp.include_routers(handlers.router) dp.include_routers(handlers.router)
await bot.delete_webhook(drop_pending_updates=True) await bot.delete_webhook(drop_pending_updates=True)
monitor = asyncio.create_task(run_monitor()) monitor = asyncio.create_task(run_monitor())
web = await run_web() if env.web.enabled else None
try: try:
await dp.start_polling(bot) await dp.start_polling(bot)
finally: finally:
monitor.cancel() monitor.cancel()
with contextlib.suppress(asyncio.CancelledError): with contextlib.suppress(asyncio.CancelledError):
await monitor await monitor
if web is not None:
await web.cleanup()
await komodo.close()
await bot.session.close() await bot.session.close()
await db.close() await db.close()
+106
View File
@@ -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"Сервер <b>{html.escape(name)}</b> не отвечает", [
f"Ошибка: {html.escape(str(d.get('err') or 'нет деталей'))}"
]
if kind == "ServerCpu":
return f"CPU на <b>{html.escape(name)}</b>", [
f"Загрузка: {_pct(d.get('percentage'))}"
]
if kind == "ServerMem":
return f"Память на <b>{html.escape(name)}</b>", [
f"Занято: {_gb(d.get('used_gb'))} из {_gb(d.get('total_gb'))}"
]
if kind == "ServerDisk":
return f"Диск на <b>{html.escape(name)}</b>", [
f"Точка: <code>{html.escape(str(d.get('path', '/')))}</code>",
f"Занято: {_gb(d.get('used_gb'))} из {_gb(d.get('total_gb'))}",
]
if kind == "ServerVersionMismatch":
return f"Разъехались версии на <b>{html.escape(name)}</b>", [
f"Periphery: {d.get('version')}, Core: {d.get('core_version')}",
"Чинится только по ssh - дашборд к такому серверу не достучится.",
]
if kind in {"StackStateChange", "ContainerStateChange"}:
return f"<b>{html.escape(name)}</b>: {d.get('from')}{d.get('to')}", (
[f"Сервер: {html.escape(server)}"] if server else []
)
if kind in {"StackImageUpdateAvailable", "DeploymentImageUpdateAvailable"}:
return f"Обновление образа для <b>{html.escape(name)}</b>", [
f"<code>{html.escape(str(d.get('image', '?')))}</code>"
]
if kind in {"StackAutoUpdated", "DeploymentAutoUpdated"}:
return f"<b>{html.escape(name)}</b> обновлён автоматически", [
f"<code>{html.escape(str(d.get('images') or d.get('image') or ''))}</code>"
]
if kind == "ResourceSyncPendingUpdates":
return f"Синк <b>{html.escape(name)}</b>: git разошёлся с Komodo", [
"Кто-то правил ресурсы мимо git, либо пуш не доехал."
]
if kind in {"ProcedureFailed", "ActionFailed", "BuildFailed", "RepoBuildFailed"}:
what = {"ProcedureFailed": "Процедура", "ActionFailed": "Действие"}.get(
kind, "Сборка"
)
return f"{what} <b>{html.escape(name)}</b> упала", []
if kind == "ScheduleRun":
return f"По расписанию запущено: <b>{html.escape(name)}</b>", []
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)}: <b>{html.escape(name)}</b>", 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"<i>Komodo · {html.escape(kind)}</i>"]
)
lines = [f"{icon} {title}"]
lines.extend(x for x in details if x)
lines.append(f"<i>Komodo · {html.escape(kind)}</i>")
return "\n".join(lines)
+111 -18
View File
@@ -1,39 +1,76 @@
"""HTTP-пробинг сервисов и логика сравнения ответа с эталоном.""" """HTTP-пробинг сервисов и логика сравнения ответа с эталоном."""
import difflib
import html
import json import json
from dataclasses import dataclass import time
from dataclasses import dataclass, field
import aiohttp import aiohttp
from utils.db.models import MatchType, Service from utils.db.models import MatchType, Service
from utils.env import env
from utils.format import short from utils.format import short
_MISSING = object() _MISSING = object()
DIFF_LINES = 12
DIFF_HEADER_LINES = 2
@dataclass(slots=True) @dataclass(slots=True)
class Probe: class Probe:
reachable: bool # удалось ли вообще получить ответ reachable: bool
status: int | None status: int | None
body: str body: str
error: str | None 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"Ожидали: <code>{html.escape(self.expected or '-')}</code>")
lines.append(f"Сейчас: <code>{html.escape(self.actual or '-')}</code>")
lines.extend(self.details)
return "\n".join(lines)
async def probe(session: aiohttp.ClientSession, url: str, timeout: int) -> Probe: # noqa: ASYNC109 async def probe(session: aiohttp.ClientSession, url: str, timeout: int) -> Probe: # noqa: ASYNC109
started = time.monotonic()
try: try:
async with session.get( async with session.get(
url, timeout=aiohttp.ClientTimeout(total=timeout), allow_redirects=True url, timeout=aiohttp.ClientTimeout(total=timeout), allow_redirects=True
) as resp: ) as resp:
body = await resp.text() 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: 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: except aiohttp.ClientError as exc:
return Probe(reachable=False, status=None, body="", error=f"нет связи: {exc}") return Probe(reachable=False, status=None, body="", error=f"нет связи: {exc}")
def dig(data: object, path: str) -> object: def dig(data: object, path: str) -> object:
"""Достаёт значение по dot-пути (индексы списков числами).""" """Достаёт значение по dot-пути (индексы списков - числами)."""
node = data node = data
for segment in path.split("."): for segment in path.split("."):
if isinstance(node, dict): if isinstance(node, dict):
@@ -49,32 +86,88 @@ def dig(data: object, path: str) -> object:
return node return node
def evaluate(service: Service, result: Probe) -> tuple[bool, str]: # noqa: PLR0911 def snippet(text: str, limit: int | None = None) -> list[str]:
"""Возвращает (жив ли сервис, причина падения).""" """Кусок ответа в <pre>, готовый к вставке в сообщение."""
limit = limit or env.monitor.body_preview
body = text.strip()
if not body:
return ["<i>(пустой ответ)</i>"]
clipped = body[:limit]
tail = "\n" if len(body) > limit else ""
return [f"<pre>{html.escape(clipped)}{tail}</pre>"]
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"<pre>{html.escape(text)}{more}</pre>"]
def evaluate(service: Service, result: Probe) -> Verdict: # noqa: PLR0911
"""Жив ли сервис - и если нет, то что именно он сейчас отдаёт."""
if not result.reachable: 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 service.match_type is MatchType.STATUS:
if result.status != service.ok_status: if result.status != service.ok_status:
return False, f"HTTP {result.status} (ждали {service.ok_status})" return Verdict(
return True, "" 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 service.match_type is MatchType.BODY:
if result.body.strip() != service.ok_body.strip(): if result.body.strip() != service.ok_body.strip():
return False, "тело ответа изменилось" return Verdict(
return True, "" healthy=False,
reason="тело ответа изменилось",
details=body_diff(service.ok_body, result.body),
)
return Verdict(healthy=True)
# MatchType.JSON
try: try:
data = json.loads(result.body) data = json.loads(result.body)
except ValueError: 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 "") value = dig(data, service.json_path or "")
if value is _MISSING: if value is _MISSING:
return False, f"нет поля {service.json_path}" return Verdict(
if _scalar_str(value) != service.json_value: healthy=False,
return False, f"{service.json_path} = {short(_scalar_str(value))}" reason=f"нет поля {service.json_path}",
return True, "" 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: def _scalar_str(value: object) -> str:
+3 -1
View File
@@ -5,8 +5,10 @@ from utils.env import env
class Admin(BaseFilter): class Admin(BaseFilter):
"""Пускаем только своих."""
async def __call__(self, event: TelegramObject) -> bool: async def __call__(self, event: TelegramObject) -> bool:
user = None user = None
if isinstance(event, (Message, CallbackQuery)): if isinstance(event, (Message, CallbackQuery)):
user = event.from_user 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
+2 -2
View File
@@ -1,8 +1,8 @@
from aiogram import Router from aiogram import Router
from . import add, commands from . import add, commands, fleet
router = Router() router = Router()
router.include_routers(commands.router, add.router) router.include_routers(commands.router, add.router, fleet.router)
__all__ = ["router"] __all__ = ["router"]
+4 -4
View File
@@ -83,7 +83,7 @@ def _choose_view(data: dict) -> tuple[str, types.InlineKeyboardMarkup]:
if len(kids) > MAX_CHILDREN: if len(kids) > MAX_CHILDREN:
lines.append(f"…показаны первые {MAX_CHILDREN} полей") lines.append(f"…показаны первые {MAX_CHILDREN} полей")
else: else:
lines.append("Ответ не JSON можно следить за статусом или телом целиком.") lines.append("Ответ не JSON - можно следить за статусом или телом целиком.")
if path: if path:
kb.button(text="⬆️ Назад", callback_data="add:up") 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): if isinstance(callback.message, types.Message):
await callback.message.edit_text( await callback.message.edit_text(
"Пришли ссылку на сервис " "Пришли ссылку на сервис "
"(например <code>https://example.com/health</code>) " "(например <code>https://example.com/health</code>) - "
"я запрошу её и предложу выбрать, за чем следить.", "я запрошу её и предложу выбрать, за чем следить.",
reply_markup=back_to_menu(), reply_markup=back_to_menu(),
) )
@@ -112,7 +112,7 @@ async def on_url(message: types.Message, state: FSMContext) -> None:
url = (message.text or "").strip() url = (message.text or "").strip()
if await db.services.by_url(url) is not None: if await db.services.by_url(url) is not None:
await message.answer("Этот сервис уже под наблюдением. /list список.") await message.answer("Этот сервис уже под наблюдением. /list - список.")
return return
note = await message.answer("Запрашиваю…") note = await message.answer("Запрашиваю…")
@@ -123,7 +123,7 @@ async def on_url(message: types.Message, state: FSMContext) -> None:
await state.clear() await state.clear()
await note.edit_text( await note.edit_text(
f"Не смог достучаться: <b>{html.escape(result.error or '')}</b>.\n" f"Не смог достучаться: <b>{html.escape(result.error or '')}</b>.\n"
"Эталон снимаю только с живого сервиса попробуй, когда он поднимется.", "Эталон снимаю только с живого сервиса - попробуй, когда он поднимется.",
reply_markup=back_to_menu(), reply_markup=back_to_menu(),
) )
return return
+193 -20
View File
@@ -1,12 +1,18 @@
"""Меню, список сервисов, статистика, удаление.""" """Меню, список сервисов, карточка, статистика, ручная проверка, удаление."""
import html import html
import json as jsonlib
import aiohttp
from aiogram import Bot, F, Router, types from aiogram import Bot, F, Router, types
from aiogram.filters import Command, CommandStart 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.filters import Admin
from healthbot.keyboards import ( from healthbot.keyboards import (
after_check,
back_to_menu, back_to_menu,
confirm_delete, confirm_delete,
main_menu, main_menu,
@@ -14,7 +20,8 @@ from healthbot.keyboards import (
services_list, services_list,
) )
from utils.db import db 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.format import human_duration, now, parse
from utils.logging import logger from utils.logging import logger
@@ -22,21 +29,28 @@ router = Router()
router.message.filter(Admin()) router.message.filter(Admin())
router.callback_query.filter(Admin()) router.callback_query.filter(Admin())
MAX_NAME = 64
GREETING = ( GREETING = (
"👋 <b>healthbot</b> слежу за твоими сервисами и пишу, если что-то легло.\n\n" "👋 <b>healthbot</b> - слежу за твоими сервисами и пишу, если что-то легло.\n\n"
"Пришли ссылку на health-эндпоинт или жми кнопку 👇" "Пришли ссылку на health-эндпоинт или жми кнопку 👇"
) )
class RenameFlow(StatesGroup):
waiting_name = State()
@router.startup() @router.startup()
async def on_startup(bot: Bot) -> None: async def on_startup(bot: Bot) -> None:
await bot.set_my_commands( commands = [
[ types.BotCommand(command="start", description="Меню"),
types.BotCommand(command="start", description="Меню"), types.BotCommand(command="list", description="Сервисы"),
types.BotCommand(command="list", description="Сервисы"), types.BotCommand(command="stats", 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}") 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() services = await db.services.all()
if not services: if not services:
return "Пока ни одного сервиса. Пришли ссылку, чтобы добавить.", back_to_menu() return "Пока ни одного сервиса. Пришли ссылку, чтобы добавить.", back_to_menu()
up = sum(s.is_up for s in services) watched = [s for s in services if not s.is_paused]
text = f"📋 Сервисы — 🟢 {up} / 🔴 {len(services) - up}\nТыкни, чтобы открыть." 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) return text, services_list(services)
@@ -94,8 +111,15 @@ async def _card_view(service: Service) -> str:
window = (now() - parse(service.created_at)).total_seconds() window = (now() - parse(service.created_at)).total_seconds()
uptime = 100.0 if window <= 0 else max(0.0, (window - down) / window * 100) 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 = "🟢 работает" status = "🟢 работает"
if service.fail_streak:
status += (
f"\n ⚠️ неудач подряд: {service.fail_streak}"
f"/{env.monitor.failure_threshold}"
)
else: else:
current = await db.incidents.current(service.id) current = await db.incidents.current(service.id)
since = ( since = (
@@ -106,31 +130,176 @@ async def _card_view(service: Service) -> str:
reason = html.escape(current.reason) if current else "" reason = html.escape(current.reason) if current else ""
status = f"🔴 лежит уже {since}\n причина: {reason}" status = f"🔴 лежит уже {since}\n причина: {reason}"
latency = (
f"{service.last_latency_ms} мс" if service.last_latency_ms is not None else "-"
)
return ( return (
f"<b>{html.escape(service.name)}</b>\n" f"<b>{html.escape(service.name)}</b>\n"
f"<code>{html.escape(service.url)}</code>\n\n" f"<code>{html.escape(service.url)}</code>\n\n"
f"Правило: {service.rule}\n" f"Правило: {service.rule}\n"
f"Статус: {status}\n\n" f"Статус: {status}\n"
f"Отклик: {latency}\n\n"
f"📊 Аптайм: <b>{uptime:.2f}%</b>\n" f"📊 Аптайм: <b>{uptime:.2f}%</b>\n"
f"Инцидентов: {len(incidents)}\n" f"Инцидентов: {len(incidents)}\n"
f"Суммарный даунтайм: {human_duration(down) if down else '0с'}" f"Суммарный даунтайм: {human_duration(down) if down else '0с'}"
) )
@router.callback_query(F.data.startswith("svc:open:")) async def _show_card(callback: types.CallbackQuery, service: Service) -> None:
async def svc_open(callback: types.CallbackQuery) -> 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_id = int((callback.data or "").rsplit(":", 1)[-1])
service = await db.services.get(service_id) service = await db.services.get(service_id)
if service is None: if service is None:
await callback.answer("Сервис не найден", show_alert=True) 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 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"<b>{html.escape(service.name)}</b>",
f"<code>{html.escape(service.url)}</code>",
"",
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): if isinstance(callback.message, types.Message):
await callback.message.edit_text( 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"Как назвать вместо <b>{html.escape(service.name)}</b>?\n"
"Пришли новое имя сообщением.",
reply_markup=back_to_menu(),
) )
await callback.answer() 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:")) @router.callback_query(F.data.startswith("svc:del:"))
async def svc_del(callback: types.CallbackQuery) -> None: async def svc_del(callback: types.CallbackQuery) -> None:
service_id = int((callback.data or "").rsplit(":", 1)[-1]) service_id = int((callback.data or "").rsplit(":", 1)[-1])
@@ -162,11 +331,15 @@ async def _stats_view() -> tuple[str, types.InlineKeyboardMarkup]:
down = _downtime(incidents) down = _downtime(incidents)
window = (now() - parse(service.created_at)).total_seconds() window = (now() - parse(service.created_at)).total_seconds()
uptime = 100.0 if window <= 0 else max(0.0, (window - down) / window * 100) 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( lines.append(
f"{mark} <b>{html.escape(service.name)}</b> " f"{service.mark} <b>{html.escape(service.name)}</b> - "
f"аптайм {uptime:.2f}%, инцидентов {len(incidents)}, " 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() return "\n".join(lines), back_to_menu()
+205
View File
@@ -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 = ["🛰 <b>Флот</b>", ""]
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} <b>{html.escape(server.name)}</b> - "
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} <b>{html.escape(stack.name)}</b>\n"
f"Состояние: <code>{html.escape(stack.state or '?')}</code>\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<pre>{html.escape(tail)}</pre>",
reply_markup=kb.as_markup(),
)
+27 -6
View File
@@ -2,6 +2,7 @@ from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup
from aiogram.utils.keyboard import InlineKeyboardBuilder from aiogram.utils.keyboard import InlineKeyboardBuilder
from utils.db.models import Service from utils.db.models import Service
from utils.env import env
def main_menu() -> InlineKeyboardMarkup: 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:add")
kb.button(text="📋 Сервисы", callback_data="menu:list") kb.button(text="📋 Сервисы", callback_data="menu:list")
kb.button(text="📊 Статистика", callback_data="menu:stats") 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() return kb.as_markup()
def services_list(services: list[Service]) -> InlineKeyboardMarkup: def services_list(services: list[Service]) -> InlineKeyboardMarkup:
kb = InlineKeyboardBuilder() kb = InlineKeyboardBuilder()
for svc in services: for svc in services:
mark = "🟢" if svc.is_up else "🔴" kb.button(text=f"{svc.mark} {svc.name}", callback_data=f"svc:open:{svc.id}")
kb.button(text=f"{mark} {svc.name}", callback_data=f"svc:open:{svc.id}")
kb.button(text=" Добавить", callback_data="menu:add") kb.button(text=" Добавить", callback_data="menu:add")
kb.button(text="⬅️ Меню", callback_data="menu:main") kb.button(text="⬅️ Меню", callback_data="menu:main")
kb.adjust(1) kb.adjust(1)
return kb.as_markup() return kb.as_markup()
def service_card(service_id: int) -> InlineKeyboardMarkup: def service_card(service: Service) -> InlineKeyboardMarkup:
kb = InlineKeyboardBuilder() 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.button(text="⬅️ К списку", callback_data="menu:list")
kb.adjust(2) kb.adjust(1, 2, 2)
return kb.as_markup() return kb.as_markup()
@@ -40,6 +50,17 @@ def confirm_delete(service_id: int) -> InlineKeyboardMarkup:
return kb.as_markup() 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: def back_to_menu() -> InlineKeyboardMarkup:
return InlineKeyboardMarkup( return InlineKeyboardMarkup(
inline_keyboard=[ inline_keyboard=[
+132
View File
@@ -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()
+77 -43
View File
@@ -5,8 +5,8 @@ import html
import aiohttp import aiohttp
from healthbot.checker import evaluate, probe from healthbot.checker import Verdict, evaluate, probe
from healthbot.common import bot from healthbot.notify import notify
from utils.db import db from utils.db import db
from utils.db.models import Service from utils.db.models import Service
from utils.env import env from utils.env import env
@@ -14,62 +14,96 @@ from utils.format import human_duration, now, parse
from utils.logging import logger from utils.logging import logger
async def _notify(text: str) -> None: def _head(service: Service, icon: str, tail: str) -> str:
try: return (
await bot.send_message(env.admin_id, text) f"{icon} <b>{html.escape(service.name)}</b> {tail}\n"
except Exception: f"<code>{html.escape(service.url)}</code>"
logger.exception("Не смог отправить уведомление админу") )
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) result = await probe(session, service.url, env.monitor.request_timeout)
healthy, reason = evaluate(service, result) verdict = evaluate(service, result)
name = html.escape(service.name) await db.services.touch_probe(service.id, result.latency_ms)
if healthy and not service.is_up: if verdict.healthy:
current = await db.incidents.current(service.id) if service.fail_streak:
downtime = "" await db.services.set_streak(service.id, 0)
if current is not None: if not service.is_up:
await db.incidents.close(current.id) await _on_recovered(service)
downtime = human_duration( return verdict
(now() - parse(current.started_at)).total_seconds()
)
await db.services.set_up(service.id, is_up=True)
await _notify(f"✅ <b>{name}</b> снова в строю.\nДаунтайм: {downtime}.")
logger.info(f"[green]UP[/] {service.name}")
elif not healthy and service.is_up: streak = service.fail_streak + 1
await db.incidents.open(service.id, reason) await db.services.set_streak(service.id, streak)
await db.services.set_up(service.id, is_up=False)
await _notify(
f"🔴 <b>{name}</b> лёг.\n"
f"<code>{html.escape(service.url)}</code>\n"
f"Причина: {html.escape(reason)}"
)
logger.warning(f"[red]DOWN[/] {service.name}: {reason}")
elif not healthy and not service.is_up: if service.is_up:
current = await db.incidents.current(service.id) if streak < env.monitor.failure_threshold:
if current is None: logger.info(
return f"[yellow]FLAP[/] {service.name}: "
idle = (now() - parse(current.last_reminder_at)).total_seconds() f"{streak}/{env.monitor.failure_threshold} - {verdict.reason}"
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"⏰ <b>{name}</b> всё ещё лежит (уже {since}).\n"
f"Причина: {html.escape(reason)}"
) )
return verdict
await _on_failed(service, verdict)
else:
await _on_still_down(service, verdict)
return verdict
async def run_monitor() -> None: 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: async with aiohttp.ClientSession() as session:
while True: while True:
try: try:
for service in await db.services.all(): for service in await db.services.all():
if service.is_paused:
continue
try: try:
await _check_one(session, service) await check_one(session, service)
except Exception: except Exception:
logger.exception(f"Ошибка проверки {service.url}") logger.exception(f"Ошибка проверки {service.url}")
except Exception: except Exception:
+28
View File
@@ -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("<pre>") > cut.count("</pre>"):
cut += "</pre>"
return cut + "\n<i>…обрезано</i>"
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}")
+69
View File
@@ -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<pre>{raw[:1000]}</pre>"
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
+17
View File
@@ -32,6 +32,13 @@ CREATE TABLE IF NOT EXISTS incidents (
CREATE INDEX IF NOT EXISTS idx_incidents_service ON incidents(service_id); 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: class Database:
def __init__(self, path: str) -> None: def __init__(self, path: str) -> None:
@@ -46,10 +53,20 @@ class Database:
self._conn.row_factory = aiosqlite.Row self._conn.row_factory = aiosqlite.Row
await self._conn.execute("PRAGMA foreign_keys = ON") await self._conn.execute("PRAGMA foreign_keys = ON")
await self._conn.executescript(SCHEMA) await self._conn.executescript(SCHEMA)
await self._migrate()
await self._conn.commit() await self._conn.commit()
self.services = ServiceRepository(self._conn) self.services = ServiceRepository(self._conn)
self.incidents = IncidentRepository(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: async def close(self) -> None:
if self._conn is not None: if self._conn is not None:
await self._conn.close() await self._conn.close()
+13 -3
View File
@@ -3,9 +3,9 @@ from enum import StrEnum
class MatchType(StrEnum): class MatchType(StrEnum):
STATUS = "status" # следим только за HTTP-статусом STATUS = "status"
BODY = "body" # тело ответа должно совпадать целиком BODY = "body"
JSON = "json" # конкретное поле JSON должно быть равно эталону JSON = "json"
@dataclass(slots=True) @dataclass(slots=True)
@@ -20,6 +20,10 @@ class Service:
json_value: str | None json_value: str | None
is_up: bool is_up: bool
created_at: str created_at: str
fail_streak: int = 0
is_paused: bool = False
last_latency_ms: int | None = None
last_checked_at: str | None = None
@property @property
def rule(self) -> str: def rule(self) -> str:
@@ -29,6 +33,12 @@ class Service:
return f"<code>{self.json_path}</code> = <code>{self.json_value}</code>" return f"<code>{self.json_path}</code> = <code>{self.json_value}</code>"
return "тело ответа неизменно" return "тело ответа неизменно"
@property
def mark(self) -> str:
if self.is_paused:
return ""
return "🟢" if self.is_up else "🔴"
@dataclass(slots=True) @dataclass(slots=True)
class Incident: class Incident:
+43 -1
View File
@@ -5,7 +5,8 @@ from utils.format import now
_COLUMNS = ( _COLUMNS = (
"id, name, url, match_type, ok_status, ok_body, " "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"], json_value=row["json_value"],
is_up=bool(row["is_up"]), is_up=bool(row["is_up"]),
created_at=row["created_at"], 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() 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: async def delete(self, service_id: int) -> None:
await self._conn.execute("DELETE FROM services WHERE id = ?", (service_id,)) await self._conn.execute("DELETE FROM services WHERE id = ?", (service_id,))
await self._conn.commit() await self._conn.commit()
+41 -5
View File
@@ -1,4 +1,4 @@
from pydantic import Field, SecretStr from pydantic import Field, SecretStr, computed_field
from pydantic_settings import BaseSettings, SettingsConfigDict from pydantic_settings import BaseSettings, SettingsConfigDict
@@ -7,9 +7,36 @@ class BotSettings(BaseSettings):
class MonitorSettings(BaseSettings): class MonitorSettings(BaseSettings):
check_interval: int = 30 # как часто пинговать сервисы, сек check_interval: int = 30
reminder_interval: int = 1800 # как часто напоминать что всё ещё лежит, сек reminder_interval: int = 1800
request_timeout: int = 10 # таймаут одного запроса, сек 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): class DbSettings(BaseSettings):
@@ -24,9 +51,13 @@ class LogSettings(BaseSettings):
class Settings(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) bot: BotSettings = Field(default_factory=BotSettings)
monitor: MonitorSettings = Field(default_factory=MonitorSettings) monitor: MonitorSettings = Field(default_factory=MonitorSettings)
komodo: KomodoSettings = Field(default_factory=KomodoSettings)
web: WebSettings = Field(default_factory=WebSettings)
db: DbSettings = Field(default_factory=DbSettings) db: DbSettings = Field(default_factory=DbSettings)
log: LogSettings = Field(default_factory=LogSettings) log: LogSettings = Field(default_factory=LogSettings)
@@ -34,5 +65,10 @@ class Settings(BaseSettings):
case_sensitive=False, env_file=".env", env_nested_delimiter="__", extra="ignore" 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() env = Settings()