feat(bot): please claude i need this my services are kinda homeless
This commit is contained in:
@@ -0,0 +1,37 @@
|
||||
import asyncio
|
||||
import contextlib
|
||||
|
||||
from utils.logging import logger, setup_logging
|
||||
|
||||
setup_logging()
|
||||
|
||||
|
||||
async def runner() -> None:
|
||||
from utils.db import db # noqa: PLC0415
|
||||
|
||||
from . import handlers # noqa: PLC0415
|
||||
from .common import bot, dp # noqa: PLC0415
|
||||
from .monitor import run_monitor # noqa: PLC0415
|
||||
|
||||
await db.connect()
|
||||
dp.include_routers(handlers.router)
|
||||
await bot.delete_webhook(drop_pending_updates=True)
|
||||
|
||||
monitor = asyncio.create_task(run_monitor())
|
||||
try:
|
||||
await dp.start_polling(bot)
|
||||
finally:
|
||||
monitor.cancel()
|
||||
with contextlib.suppress(asyncio.CancelledError):
|
||||
await monitor
|
||||
await bot.session.close()
|
||||
await db.close()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
logger.info("Starting...")
|
||||
|
||||
with contextlib.suppress(KeyboardInterrupt):
|
||||
asyncio.run(runner())
|
||||
|
||||
logger.info("[red]Stopped.[/]")
|
||||
@@ -0,0 +1,4 @@
|
||||
from . import main
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,86 @@
|
||||
"""HTTP-пробинг сервисов и логика сравнения ответа с эталоном."""
|
||||
|
||||
import json
|
||||
from dataclasses import dataclass
|
||||
|
||||
import aiohttp
|
||||
|
||||
from utils.db.models import MatchType, Service
|
||||
from utils.format import short
|
||||
|
||||
_MISSING = object()
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class Probe:
|
||||
reachable: bool # удалось ли вообще получить ответ
|
||||
status: int | None
|
||||
body: str
|
||||
error: str | None
|
||||
|
||||
|
||||
async def probe(session: aiohttp.ClientSession, url: str, timeout: int) -> Probe: # noqa: ASYNC109
|
||||
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)
|
||||
except TimeoutError:
|
||||
return Probe(reachable=False, status=None, body="", error="таймаут")
|
||||
except aiohttp.ClientError as exc:
|
||||
return Probe(reachable=False, status=None, body="", error=f"нет связи: {exc}")
|
||||
|
||||
|
||||
def dig(data: object, path: str) -> object:
|
||||
"""Достаёт значение по dot-пути (индексы списков — числами)."""
|
||||
node = data
|
||||
for segment in path.split("."):
|
||||
if isinstance(node, dict):
|
||||
if segment not in node:
|
||||
return _MISSING
|
||||
node = node[segment]
|
||||
elif isinstance(node, list):
|
||||
if not segment.isdigit() or int(segment) >= len(node):
|
||||
return _MISSING
|
||||
node = node[int(segment)]
|
||||
else:
|
||||
return _MISSING
|
||||
return node
|
||||
|
||||
|
||||
def evaluate(service: Service, result: Probe) -> tuple[bool, str]: # noqa: PLR0911
|
||||
"""Возвращает (жив ли сервис, причина падения)."""
|
||||
if not result.reachable:
|
||||
return False, 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, ""
|
||||
|
||||
if service.match_type is MatchType.BODY:
|
||||
if result.body.strip() != service.ok_body.strip():
|
||||
return False, "тело ответа изменилось"
|
||||
return True, ""
|
||||
|
||||
# MatchType.JSON
|
||||
try:
|
||||
data = json.loads(result.body)
|
||||
except ValueError:
|
||||
return False, "ответ перестал быть JSON"
|
||||
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, ""
|
||||
|
||||
|
||||
def _scalar_str(value: object) -> str:
|
||||
"""Строковое представление скаляра JSON (bool/null как в JSON)."""
|
||||
if isinstance(value, bool):
|
||||
return "true" if value else "false"
|
||||
if value is None:
|
||||
return "null"
|
||||
return str(value)
|
||||
@@ -0,0 +1,13 @@
|
||||
from aiogram import Bot, Dispatcher
|
||||
from aiogram.client.default import DefaultBotProperties
|
||||
from aiogram.fsm.storage.memory import MemoryStorage
|
||||
|
||||
from utils.env import env
|
||||
|
||||
bot = Bot(
|
||||
token=env.bot.token.get_secret_value(),
|
||||
default=DefaultBotProperties(parse_mode="HTML"),
|
||||
)
|
||||
dp = Dispatcher(storage=MemoryStorage())
|
||||
|
||||
__all__ = ["bot", "dp"]
|
||||
@@ -0,0 +1,3 @@
|
||||
from .admin import Admin
|
||||
|
||||
__all__ = ["Admin"]
|
||||
@@ -0,0 +1,12 @@
|
||||
from aiogram.filters import BaseFilter
|
||||
from aiogram.types import CallbackQuery, Message, TelegramObject
|
||||
|
||||
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
|
||||
@@ -0,0 +1,8 @@
|
||||
from aiogram import Router
|
||||
|
||||
from . import add, commands
|
||||
|
||||
router = Router()
|
||||
router.include_routers(commands.router, add.router)
|
||||
|
||||
__all__ = ["router"]
|
||||
@@ -0,0 +1,240 @@
|
||||
"""Добавление сервиса: прислал ссылку → выбрал в меню за чем следить."""
|
||||
|
||||
import html
|
||||
from urllib.parse import urlparse
|
||||
|
||||
import aiohttp
|
||||
from aiogram import F, Router, types
|
||||
from aiogram.fsm.context import FSMContext
|
||||
from aiogram.fsm.state import State, StatesGroup
|
||||
from aiogram.utils.keyboard import InlineKeyboardBuilder
|
||||
|
||||
from healthbot.checker import _scalar_str, probe
|
||||
from healthbot.filters import Admin
|
||||
from healthbot.keyboards import back_to_menu
|
||||
from utils.db import db
|
||||
from utils.db.models import MatchType
|
||||
from utils.env import env
|
||||
|
||||
router = Router()
|
||||
router.message.filter(Admin())
|
||||
router.callback_query.filter(Admin())
|
||||
|
||||
URL_RE = r"^https?://\S+$"
|
||||
MAX_CHILDREN = 40
|
||||
PREVIEW_LEN = 30
|
||||
|
||||
|
||||
class AddFlow(StatesGroup):
|
||||
waiting_url = State()
|
||||
choosing = State()
|
||||
|
||||
|
||||
def _name_from_url(url: str) -> str:
|
||||
parsed = urlparse(url)
|
||||
name = parsed.netloc + (parsed.path if parsed.path not in ("", "/") else "")
|
||||
return name[:60] or url[:60]
|
||||
|
||||
|
||||
def _node_at(root: object, path: list[str]) -> object:
|
||||
node = root
|
||||
for seg in path:
|
||||
node = node[int(seg)] if isinstance(node, list) else node[seg] # type: ignore[index]
|
||||
return node
|
||||
|
||||
|
||||
def _children(node: object) -> list[tuple[str, object]]:
|
||||
if isinstance(node, dict):
|
||||
return list(node.items())
|
||||
if isinstance(node, list):
|
||||
return [(str(i), v) for i, v in enumerate(node)]
|
||||
return []
|
||||
|
||||
|
||||
def _preview(value: object) -> str:
|
||||
if isinstance(value, dict):
|
||||
return "{…}"
|
||||
if isinstance(value, list):
|
||||
return "[…]"
|
||||
text = _scalar_str(value)
|
||||
return text if len(text) <= PREVIEW_LEN else text[: PREVIEW_LEN - 1] + "…"
|
||||
|
||||
|
||||
def _choose_view(data: dict) -> tuple[str, types.InlineKeyboardMarkup]:
|
||||
parsed = data["parsed"]
|
||||
path: list[str] = data["path"]
|
||||
body_preview = html.escape(data["body"][:350]) or "(пусто)"
|
||||
|
||||
lines = [
|
||||
f"🔗 <b>{html.escape(data['name'])}</b>",
|
||||
f"<code>{html.escape(data['url'])}</code>",
|
||||
f"Ответ: HTTP {data['status']}",
|
||||
f"<pre>{body_preview}</pre>",
|
||||
]
|
||||
|
||||
kb = InlineKeyboardBuilder()
|
||||
if parsed is not None:
|
||||
crumb = " › ".join(["корень", *path])
|
||||
lines.append(f"За каким полем следить? <i>{html.escape(crumb)}</i>")
|
||||
node = _node_at(parsed, path)
|
||||
kids = _children(node)
|
||||
for i, (key, value) in enumerate(kids[:MAX_CHILDREN]):
|
||||
kb.button(text=f"{key}: {_preview(value)}", callback_data=f"add:nav:{i}")
|
||||
if len(kids) > MAX_CHILDREN:
|
||||
lines.append(f"…показаны первые {MAX_CHILDREN} полей")
|
||||
else:
|
||||
lines.append("Ответ не JSON — можно следить за статусом или телом целиком.")
|
||||
|
||||
if path:
|
||||
kb.button(text="⬆️ Назад", callback_data="add:up")
|
||||
kb.button(text="🔢 Только HTTP-статус", callback_data="add:status")
|
||||
kb.button(text="📄 Всё тело целиком", callback_data="add:body")
|
||||
kb.button(text="❌ Отмена", callback_data="add:cancel")
|
||||
kb.adjust(1)
|
||||
return "\n".join(lines), kb.as_markup()
|
||||
|
||||
|
||||
@router.callback_query(F.data == "menu:add")
|
||||
async def on_add_menu(callback: types.CallbackQuery, state: FSMContext) -> None:
|
||||
await state.set_state(AddFlow.waiting_url)
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text(
|
||||
"Пришли ссылку на сервис "
|
||||
"(например <code>https://example.com/health</code>) — "
|
||||
"я запрошу её и предложу выбрать, за чем следить.",
|
||||
reply_markup=back_to_menu(),
|
||||
)
|
||||
await callback.answer()
|
||||
|
||||
|
||||
@router.message(F.text.regexp(URL_RE))
|
||||
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 — список.")
|
||||
return
|
||||
|
||||
note = await message.answer("Запрашиваю…")
|
||||
async with aiohttp.ClientSession() as session:
|
||||
result = await probe(session, url, env.monitor.request_timeout)
|
||||
|
||||
if not result.reachable:
|
||||
await state.clear()
|
||||
await note.edit_text(
|
||||
f"Не смог достучаться: <b>{html.escape(result.error or '')}</b>.\n"
|
||||
"Эталон снимаю только с живого сервиса — попробуй, когда он поднимется.",
|
||||
reply_markup=back_to_menu(),
|
||||
)
|
||||
return
|
||||
|
||||
try:
|
||||
import json # noqa: PLC0415
|
||||
|
||||
parsed = json.loads(result.body)
|
||||
except ValueError:
|
||||
parsed = None
|
||||
|
||||
await state.set_state(AddFlow.choosing)
|
||||
await state.set_data(
|
||||
{
|
||||
"url": url,
|
||||
"name": _name_from_url(url),
|
||||
"status": result.status,
|
||||
"body": result.body,
|
||||
"parsed": parsed,
|
||||
"path": [],
|
||||
}
|
||||
)
|
||||
text, markup = _choose_view(await state.get_data())
|
||||
await note.edit_text(text, reply_markup=markup)
|
||||
|
||||
|
||||
async def _save( # noqa: PLR0913
|
||||
callback: types.CallbackQuery,
|
||||
state: FSMContext,
|
||||
*,
|
||||
match_type: MatchType,
|
||||
ok_body: str = "",
|
||||
json_path: str | None = None,
|
||||
json_value: str | None = None,
|
||||
) -> None:
|
||||
data = await state.get_data()
|
||||
service = await db.services.add(
|
||||
name=data["name"],
|
||||
url=data["url"],
|
||||
match_type=match_type,
|
||||
ok_status=data["status"],
|
||||
ok_body=ok_body,
|
||||
json_path=json_path,
|
||||
json_value=json_value,
|
||||
)
|
||||
await state.clear()
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text(
|
||||
f"✅ Слежу за <b>{html.escape(service.name)}</b>\n"
|
||||
f"<code>{html.escape(service.url)}</code>\n"
|
||||
f"Правило: {service.rule}\n"
|
||||
f"Проверяю раз в {env.monitor.check_interval}с.",
|
||||
reply_markup=back_to_menu(),
|
||||
)
|
||||
await callback.answer("Готово")
|
||||
|
||||
|
||||
@router.callback_query(AddFlow.choosing, F.data == "add:status")
|
||||
async def on_pick_status(callback: types.CallbackQuery, state: FSMContext) -> None:
|
||||
await _save(callback, state, match_type=MatchType.STATUS)
|
||||
|
||||
|
||||
@router.callback_query(AddFlow.choosing, F.data == "add:body")
|
||||
async def on_pick_body(callback: types.CallbackQuery, state: FSMContext) -> None:
|
||||
data = await state.get_data()
|
||||
await _save(callback, state, match_type=MatchType.BODY, ok_body=data["body"])
|
||||
|
||||
|
||||
@router.callback_query(AddFlow.choosing, F.data == "add:up")
|
||||
async def on_up(callback: types.CallbackQuery, state: FSMContext) -> None:
|
||||
data = await state.get_data()
|
||||
if data["path"]:
|
||||
data["path"].pop()
|
||||
await state.update_data(path=data["path"])
|
||||
await _rerender(callback, state)
|
||||
|
||||
|
||||
@router.callback_query(AddFlow.choosing, F.data.startswith("add:nav:"))
|
||||
async def on_nav(callback: types.CallbackQuery, state: FSMContext) -> None:
|
||||
data = await state.get_data()
|
||||
index = int((callback.data or "").rsplit(":", 1)[-1])
|
||||
node = _node_at(data["parsed"], data["path"])
|
||||
kids = _children(node)
|
||||
if index >= len(kids):
|
||||
await callback.answer()
|
||||
return
|
||||
key, value = kids[index]
|
||||
if isinstance(value, (dict, list)):
|
||||
data["path"].append(key)
|
||||
await state.update_data(path=data["path"])
|
||||
await _rerender(callback, state)
|
||||
return
|
||||
await _save(
|
||||
callback,
|
||||
state,
|
||||
match_type=MatchType.JSON,
|
||||
json_path=".".join([*data["path"], key]),
|
||||
json_value=_scalar_str(value),
|
||||
)
|
||||
|
||||
|
||||
@router.callback_query(AddFlow.choosing, F.data == "add:cancel")
|
||||
async def on_cancel(callback: types.CallbackQuery, state: FSMContext) -> None:
|
||||
await state.clear()
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text("Отменено.", reply_markup=back_to_menu())
|
||||
await callback.answer()
|
||||
|
||||
|
||||
async def _rerender(callback: types.CallbackQuery, state: FSMContext) -> None:
|
||||
text, markup = _choose_view(await state.get_data())
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text(text, reply_markup=markup)
|
||||
await callback.answer()
|
||||
@@ -0,0 +1,185 @@
|
||||
"""Меню, список сервисов, статистика, удаление."""
|
||||
|
||||
import html
|
||||
|
||||
from aiogram import Bot, F, Router, types
|
||||
from aiogram.filters import Command, CommandStart
|
||||
|
||||
from healthbot.filters import Admin
|
||||
from healthbot.keyboards import (
|
||||
back_to_menu,
|
||||
confirm_delete,
|
||||
main_menu,
|
||||
service_card,
|
||||
services_list,
|
||||
)
|
||||
from utils.db import db
|
||||
from utils.db.models import Service
|
||||
from utils.format import human_duration, now, parse
|
||||
from utils.logging import logger
|
||||
|
||||
router = Router()
|
||||
router.message.filter(Admin())
|
||||
router.callback_query.filter(Admin())
|
||||
|
||||
GREETING = (
|
||||
"👋 <b>healthbot</b> — слежу за твоими сервисами и пишу, если что-то легло.\n\n"
|
||||
"Пришли ссылку на health-эндпоинт или жми кнопку 👇"
|
||||
)
|
||||
|
||||
|
||||
@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="Статистика"),
|
||||
]
|
||||
)
|
||||
logger.info(f"[green]Started as[/] @{(await bot.me()).username}")
|
||||
|
||||
|
||||
@router.shutdown()
|
||||
async def on_shutdown() -> None:
|
||||
logger.info("Shutting down bot...")
|
||||
|
||||
|
||||
@router.message(CommandStart())
|
||||
async def cmd_start(message: types.Message) -> None:
|
||||
await message.answer(GREETING, reply_markup=main_menu())
|
||||
|
||||
|
||||
@router.callback_query(F.data == "menu:main")
|
||||
async def menu_main(callback: types.CallbackQuery) -> None:
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text(GREETING, reply_markup=main_menu())
|
||||
await callback.answer()
|
||||
|
||||
|
||||
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Тыкни, чтобы открыть."
|
||||
return text, services_list(services)
|
||||
|
||||
|
||||
@router.message(Command("list"))
|
||||
async def cmd_list(message: types.Message) -> None:
|
||||
text, markup = await _list_view()
|
||||
await message.answer(text, reply_markup=markup)
|
||||
|
||||
|
||||
@router.callback_query(F.data == "menu:list")
|
||||
async def menu_list(callback: types.CallbackQuery) -> None:
|
||||
text, markup = await _list_view()
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text(text, reply_markup=markup)
|
||||
await callback.answer()
|
||||
|
||||
|
||||
def _downtime(incidents: list) -> float:
|
||||
total = 0.0
|
||||
for inc in incidents:
|
||||
end = parse(inc.ended_at) if inc.ended_at else now()
|
||||
total += (end - parse(inc.started_at)).total_seconds()
|
||||
return total
|
||||
|
||||
|
||||
async def _card_view(service: Service) -> str:
|
||||
incidents = await db.incidents.for_service(service.id)
|
||||
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)
|
||||
|
||||
if service.is_up:
|
||||
status = "🟢 работает"
|
||||
else:
|
||||
current = await db.incidents.current(service.id)
|
||||
since = (
|
||||
human_duration((now() - parse(current.started_at)).total_seconds())
|
||||
if current
|
||||
else "?"
|
||||
)
|
||||
reason = html.escape(current.reason) if current else ""
|
||||
status = f"🔴 лежит уже {since}\n причина: {reason}"
|
||||
|
||||
return (
|
||||
f"<b>{html.escape(service.name)}</b>\n"
|
||||
f"<code>{html.escape(service.url)}</code>\n\n"
|
||||
f"Правило: {service.rule}\n"
|
||||
f"Статус: {status}\n\n"
|
||||
f"📊 Аптайм: <b>{uptime:.2f}%</b>\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:
|
||||
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
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text(
|
||||
await _card_view(service), reply_markup=service_card(service_id)
|
||||
)
|
||||
await callback.answer()
|
||||
|
||||
|
||||
@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])
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text(
|
||||
"Точно удалить сервис вместе со всей статистикой?",
|
||||
reply_markup=confirm_delete(service_id),
|
||||
)
|
||||
await callback.answer()
|
||||
|
||||
|
||||
@router.callback_query(F.data.startswith("svc:delyes:"))
|
||||
async def svc_del_yes(callback: types.CallbackQuery) -> None:
|
||||
service_id = int((callback.data or "").rsplit(":", 1)[-1])
|
||||
await db.services.delete(service_id)
|
||||
text, markup = await _list_view()
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text(text, reply_markup=markup)
|
||||
await callback.answer("Удалено")
|
||||
|
||||
|
||||
async def _stats_view() -> tuple[str, types.InlineKeyboardMarkup]:
|
||||
services = await db.services.all()
|
||||
if not services:
|
||||
return "Пока нет данных.", back_to_menu()
|
||||
lines = ["📊 <b>Статистика</b>\n"]
|
||||
for service in services:
|
||||
incidents = await db.incidents.for_service(service.id)
|
||||
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 "🔴"
|
||||
lines.append(
|
||||
f"{mark} <b>{html.escape(service.name)}</b> — "
|
||||
f"аптайм {uptime:.2f}%, инцидентов {len(incidents)}, "
|
||||
f"даунтайм {human_duration(down) if down else '0с'}"
|
||||
)
|
||||
return "\n".join(lines), back_to_menu()
|
||||
|
||||
|
||||
@router.message(Command("stats"))
|
||||
async def cmd_stats(message: types.Message) -> None:
|
||||
text, markup = await _stats_view()
|
||||
await message.answer(text, reply_markup=markup)
|
||||
|
||||
|
||||
@router.callback_query(F.data == "menu:stats")
|
||||
async def menu_stats(callback: types.CallbackQuery) -> None:
|
||||
text, markup = await _stats_view()
|
||||
if isinstance(callback.message, types.Message):
|
||||
await callback.message.edit_text(text, reply_markup=markup)
|
||||
await callback.answer()
|
||||
@@ -0,0 +1,48 @@
|
||||
from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup
|
||||
from aiogram.utils.keyboard import InlineKeyboardBuilder
|
||||
|
||||
from utils.db.models import Service
|
||||
|
||||
|
||||
def main_menu() -> InlineKeyboardMarkup:
|
||||
kb = InlineKeyboardBuilder()
|
||||
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)
|
||||
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="➕ Добавить", 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:
|
||||
kb = InlineKeyboardBuilder()
|
||||
kb.button(text="🗑 Удалить", callback_data=f"svc:del:{service_id}")
|
||||
kb.button(text="⬅️ К списку", callback_data="menu:list")
|
||||
kb.adjust(2)
|
||||
return kb.as_markup()
|
||||
|
||||
|
||||
def confirm_delete(service_id: int) -> InlineKeyboardMarkup:
|
||||
kb = InlineKeyboardBuilder()
|
||||
kb.button(text="✅ Да, удалить", callback_data=f"svc:delyes:{service_id}")
|
||||
kb.button(text="⬅️ Отмена", callback_data=f"svc:open:{service_id}")
|
||||
kb.adjust(2)
|
||||
return kb.as_markup()
|
||||
|
||||
|
||||
def back_to_menu() -> InlineKeyboardMarkup:
|
||||
return InlineKeyboardMarkup(
|
||||
inline_keyboard=[
|
||||
[InlineKeyboardButton(text="⬅️ Меню", callback_data="menu:main")]
|
||||
]
|
||||
)
|
||||
@@ -0,0 +1,77 @@
|
||||
"""Фоновый цикл: пингует сервисы и уведомляет о падениях/восстановлениях."""
|
||||
|
||||
import asyncio
|
||||
import html
|
||||
|
||||
import aiohttp
|
||||
|
||||
from healthbot.checker import evaluate, probe
|
||||
from healthbot.common import bot
|
||||
from utils.db import db
|
||||
from utils.db.models import Service
|
||||
from utils.env import env
|
||||
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("Не смог отправить уведомление админу")
|
||||
|
||||
|
||||
async def _check_one(session: aiohttp.ClientSession, service: Service) -> None:
|
||||
result = await probe(session, service.url, env.monitor.request_timeout)
|
||||
healthy, reason = evaluate(service, result)
|
||||
name = html.escape(service.name)
|
||||
|
||||
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"✅ <b>{name}</b> снова в строю.\nДаунтайм: {downtime}.")
|
||||
logger.info(f"[green]UP[/] {service.name}")
|
||||
|
||||
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"🔴 <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:
|
||||
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"⏰ <b>{name}</b> всё ещё лежит (уже {since}).\n"
|
||||
f"Причина: {html.escape(reason)}"
|
||||
)
|
||||
|
||||
|
||||
async def run_monitor() -> None:
|
||||
logger.info(f"Monitor: каждые {env.monitor.check_interval}с")
|
||||
async with aiohttp.ClientSession() as session:
|
||||
while True:
|
||||
try:
|
||||
for service in await db.services.all():
|
||||
try:
|
||||
await _check_one(session, service)
|
||||
except Exception:
|
||||
logger.exception(f"Ошибка проверки {service.url}")
|
||||
except Exception:
|
||||
logger.exception("Ошибка цикла мониторинга")
|
||||
await asyncio.sleep(env.monitor.check_interval)
|
||||
@@ -0,0 +1,3 @@
|
||||
from .env import env
|
||||
|
||||
__all__ = ["env"]
|
||||
@@ -0,0 +1,61 @@
|
||||
from pathlib import Path
|
||||
|
||||
import aiosqlite
|
||||
|
||||
from utils.env import env
|
||||
|
||||
from .repositories import IncidentRepository, ServiceRepository
|
||||
|
||||
SCHEMA = """
|
||||
CREATE TABLE IF NOT EXISTS services (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
name TEXT NOT NULL,
|
||||
url TEXT NOT NULL UNIQUE,
|
||||
match_type TEXT NOT NULL,
|
||||
ok_status INTEGER NOT NULL,
|
||||
ok_body TEXT NOT NULL DEFAULT '',
|
||||
json_path TEXT,
|
||||
json_value TEXT,
|
||||
is_up INTEGER NOT NULL DEFAULT 1,
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS incidents (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
service_id INTEGER NOT NULL REFERENCES services(id) ON DELETE CASCADE,
|
||||
started_at TEXT NOT NULL,
|
||||
ended_at TEXT,
|
||||
reason TEXT NOT NULL DEFAULT '',
|
||||
last_reminder_at TEXT NOT NULL
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_incidents_service ON incidents(service_id);
|
||||
"""
|
||||
|
||||
|
||||
class Database:
|
||||
def __init__(self, path: str) -> None:
|
||||
self._path = path
|
||||
self._conn: aiosqlite.Connection | None = None
|
||||
self.services: ServiceRepository = None # type: ignore[assignment]
|
||||
self.incidents: IncidentRepository = None # type: ignore[assignment]
|
||||
|
||||
async def connect(self) -> None:
|
||||
Path(self._path).parent.mkdir(parents=True, exist_ok=True)
|
||||
self._conn = await aiosqlite.connect(self._path)
|
||||
self._conn.row_factory = aiosqlite.Row
|
||||
await self._conn.execute("PRAGMA foreign_keys = ON")
|
||||
await self._conn.executescript(SCHEMA)
|
||||
await self._conn.commit()
|
||||
self.services = ServiceRepository(self._conn)
|
||||
self.incidents = IncidentRepository(self._conn)
|
||||
|
||||
async def close(self) -> None:
|
||||
if self._conn is not None:
|
||||
await self._conn.close()
|
||||
self._conn = None
|
||||
|
||||
|
||||
db = Database(env.db.path)
|
||||
|
||||
__all__ = ["Database", "db"]
|
||||
@@ -0,0 +1,40 @@
|
||||
from dataclasses import dataclass
|
||||
from enum import StrEnum
|
||||
|
||||
|
||||
class MatchType(StrEnum):
|
||||
STATUS = "status" # следим только за HTTP-статусом
|
||||
BODY = "body" # тело ответа должно совпадать целиком
|
||||
JSON = "json" # конкретное поле JSON должно быть равно эталону
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class Service:
|
||||
id: int
|
||||
name: str
|
||||
url: str
|
||||
match_type: MatchType
|
||||
ok_status: int
|
||||
ok_body: str
|
||||
json_path: str | None
|
||||
json_value: str | None
|
||||
is_up: bool
|
||||
created_at: str
|
||||
|
||||
@property
|
||||
def rule(self) -> str:
|
||||
if self.match_type is MatchType.STATUS:
|
||||
return f"HTTP-статус = {self.ok_status}"
|
||||
if self.match_type is MatchType.JSON:
|
||||
return f"<code>{self.json_path}</code> = <code>{self.json_value}</code>"
|
||||
return "тело ответа неизменно"
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class Incident:
|
||||
id: int
|
||||
service_id: int
|
||||
started_at: str
|
||||
ended_at: str | None
|
||||
reason: str
|
||||
last_reminder_at: str
|
||||
@@ -0,0 +1,4 @@
|
||||
from .incident import IncidentRepository
|
||||
from .service import ServiceRepository
|
||||
|
||||
__all__ = ["IncidentRepository", "ServiceRepository"]
|
||||
@@ -0,0 +1,74 @@
|
||||
import aiosqlite
|
||||
|
||||
from utils.db.models import Incident
|
||||
from utils.format import now
|
||||
|
||||
_COLUMNS = "id, service_id, started_at, ended_at, reason, last_reminder_at"
|
||||
|
||||
|
||||
def _row_to_incident(row: aiosqlite.Row) -> Incident:
|
||||
return Incident(
|
||||
id=row["id"],
|
||||
service_id=row["service_id"],
|
||||
started_at=row["started_at"],
|
||||
ended_at=row["ended_at"],
|
||||
reason=row["reason"],
|
||||
last_reminder_at=row["last_reminder_at"],
|
||||
)
|
||||
|
||||
|
||||
class IncidentRepository:
|
||||
def __init__(self, conn: aiosqlite.Connection) -> None:
|
||||
self._conn = conn
|
||||
|
||||
async def open(self, service_id: int, reason: str) -> Incident:
|
||||
stamp = now().isoformat()
|
||||
cursor = await self._conn.execute(
|
||||
"INSERT INTO incidents "
|
||||
"(service_id, started_at, reason, last_reminder_at) VALUES (?, ?, ?, ?)",
|
||||
(service_id, stamp, reason, stamp),
|
||||
)
|
||||
await self._conn.commit()
|
||||
incident = await self.get(cursor.lastrowid) # type: ignore[arg-type]
|
||||
assert incident is not None
|
||||
return incident
|
||||
|
||||
async def get(self, incident_id: int) -> Incident | None:
|
||||
cursor = await self._conn.execute(
|
||||
f"SELECT {_COLUMNS} FROM incidents WHERE id = ?", # noqa: S608
|
||||
(incident_id,),
|
||||
)
|
||||
row = await cursor.fetchone()
|
||||
return _row_to_incident(row) if row else None
|
||||
|
||||
async def current(self, service_id: int) -> Incident | None:
|
||||
cursor = await self._conn.execute(
|
||||
f"SELECT {_COLUMNS} FROM incidents " # noqa: S608
|
||||
"WHERE service_id = ? AND ended_at IS NULL "
|
||||
"ORDER BY id DESC LIMIT 1",
|
||||
(service_id,),
|
||||
)
|
||||
row = await cursor.fetchone()
|
||||
return _row_to_incident(row) if row else None
|
||||
|
||||
async def close(self, incident_id: int) -> None:
|
||||
await self._conn.execute(
|
||||
"UPDATE incidents SET ended_at = ? WHERE id = ?",
|
||||
(now().isoformat(), incident_id),
|
||||
)
|
||||
await self._conn.commit()
|
||||
|
||||
async def touch_reminder(self, incident_id: int) -> None:
|
||||
await self._conn.execute(
|
||||
"UPDATE incidents SET last_reminder_at = ? WHERE id = ?",
|
||||
(now().isoformat(), incident_id),
|
||||
)
|
||||
await self._conn.commit()
|
||||
|
||||
async def for_service(self, service_id: int) -> list[Incident]:
|
||||
cursor = await self._conn.execute(
|
||||
f"SELECT {_COLUMNS} FROM incidents " # noqa: S608
|
||||
"WHERE service_id = ? ORDER BY id",
|
||||
(service_id,),
|
||||
)
|
||||
return [_row_to_incident(row) for row in await cursor.fetchall()]
|
||||
@@ -0,0 +1,92 @@
|
||||
import aiosqlite
|
||||
|
||||
from utils.db.models import MatchType, Service
|
||||
from utils.format import now
|
||||
|
||||
_COLUMNS = (
|
||||
"id, name, url, match_type, ok_status, ok_body, "
|
||||
"json_path, json_value, is_up, created_at"
|
||||
)
|
||||
|
||||
|
||||
def _row_to_service(row: aiosqlite.Row) -> Service:
|
||||
return Service(
|
||||
id=row["id"],
|
||||
name=row["name"],
|
||||
url=row["url"],
|
||||
match_type=MatchType(row["match_type"]),
|
||||
ok_status=row["ok_status"],
|
||||
ok_body=row["ok_body"],
|
||||
json_path=row["json_path"],
|
||||
json_value=row["json_value"],
|
||||
is_up=bool(row["is_up"]),
|
||||
created_at=row["created_at"],
|
||||
)
|
||||
|
||||
|
||||
class ServiceRepository:
|
||||
def __init__(self, conn: aiosqlite.Connection) -> None:
|
||||
self._conn = conn
|
||||
|
||||
async def all(self) -> list[Service]:
|
||||
cursor = await self._conn.execute(
|
||||
f"SELECT {_COLUMNS} FROM services ORDER BY id" # noqa: S608
|
||||
)
|
||||
return [_row_to_service(row) for row in await cursor.fetchall()]
|
||||
|
||||
async def get(self, service_id: int) -> Service | None:
|
||||
cursor = await self._conn.execute(
|
||||
f"SELECT {_COLUMNS} FROM services WHERE id = ?", # noqa: S608
|
||||
(service_id,),
|
||||
)
|
||||
row = await cursor.fetchone()
|
||||
return _row_to_service(row) if row else None
|
||||
|
||||
async def by_url(self, url: str) -> Service | None:
|
||||
cursor = await self._conn.execute(
|
||||
f"SELECT {_COLUMNS} FROM services WHERE url = ?", # noqa: S608
|
||||
(url,),
|
||||
)
|
||||
row = await cursor.fetchone()
|
||||
return _row_to_service(row) if row else None
|
||||
|
||||
async def add( # noqa: PLR0913
|
||||
self,
|
||||
*,
|
||||
name: str,
|
||||
url: str,
|
||||
match_type: MatchType,
|
||||
ok_status: int,
|
||||
ok_body: str = "",
|
||||
json_path: str | None = None,
|
||||
json_value: str | None = None,
|
||||
) -> Service:
|
||||
cursor = await self._conn.execute(
|
||||
"INSERT INTO services "
|
||||
"(name, url, match_type, ok_status, ok_body, json_path, json_value, "
|
||||
"is_up, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, 1, ?)",
|
||||
(
|
||||
name,
|
||||
url,
|
||||
match_type.value,
|
||||
ok_status,
|
||||
ok_body,
|
||||
json_path,
|
||||
json_value,
|
||||
now().isoformat(),
|
||||
),
|
||||
)
|
||||
await self._conn.commit()
|
||||
service = await self.get(cursor.lastrowid) # type: ignore[arg-type]
|
||||
assert service is not None
|
||||
return service
|
||||
|
||||
async def set_up(self, service_id: int, *, is_up: bool) -> None:
|
||||
await self._conn.execute(
|
||||
"UPDATE services SET is_up = ? WHERE id = ?", (int(is_up), 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()
|
||||
@@ -0,0 +1,38 @@
|
||||
from pydantic import Field, SecretStr
|
||||
from pydantic_settings import BaseSettings, SettingsConfigDict
|
||||
|
||||
|
||||
class BotSettings(BaseSettings):
|
||||
token: SecretStr
|
||||
|
||||
|
||||
class MonitorSettings(BaseSettings):
|
||||
check_interval: int = 30 # как часто пинговать сервисы, сек
|
||||
reminder_interval: int = 1800 # как часто напоминать что всё ещё лежит, сек
|
||||
request_timeout: int = 10 # таймаут одного запроса, сек
|
||||
|
||||
|
||||
class DbSettings(BaseSettings):
|
||||
path: str = "data/healthbot.db"
|
||||
|
||||
|
||||
class LogSettings(BaseSettings):
|
||||
level: str = "INFO"
|
||||
level_external: str = "WARNING"
|
||||
show_time: bool = False
|
||||
console_width: int = 150
|
||||
|
||||
|
||||
class Settings(BaseSettings):
|
||||
admin_id: int
|
||||
bot: BotSettings = Field(default_factory=BotSettings)
|
||||
monitor: MonitorSettings = Field(default_factory=MonitorSettings)
|
||||
db: DbSettings = Field(default_factory=DbSettings)
|
||||
log: LogSettings = Field(default_factory=LogSettings)
|
||||
|
||||
model_config = SettingsConfigDict(
|
||||
case_sensitive=False, env_file=".env", env_nested_delimiter="__", extra="ignore"
|
||||
)
|
||||
|
||||
|
||||
env = Settings()
|
||||
@@ -0,0 +1,32 @@
|
||||
from datetime import UTC, datetime
|
||||
|
||||
|
||||
def now() -> datetime:
|
||||
return datetime.now(UTC)
|
||||
|
||||
|
||||
def parse(ts: str) -> datetime:
|
||||
return datetime.fromisoformat(ts)
|
||||
|
||||
|
||||
MINUTE = 60
|
||||
MAX_PARTS = 2
|
||||
|
||||
|
||||
def human_duration(seconds: float) -> str:
|
||||
seconds = int(seconds)
|
||||
if seconds < MINUTE:
|
||||
return f"{seconds}с"
|
||||
parts: list[str] = []
|
||||
for unit, size in (("д", 86400), ("ч", 3600), ("м", 60)):
|
||||
if seconds >= size:
|
||||
parts.append(f"{seconds // size}{unit}")
|
||||
seconds %= size
|
||||
if seconds and len(parts) < MAX_PARTS:
|
||||
parts.append(f"{seconds}с")
|
||||
return " ".join(parts[:MAX_PARTS])
|
||||
|
||||
|
||||
def short(value: object, limit: int = 60) -> str:
|
||||
text = str(value).replace("\n", " ").strip()
|
||||
return text if len(text) <= limit else text[: limit - 1] + "…"
|
||||
@@ -0,0 +1,36 @@
|
||||
import logging
|
||||
|
||||
from rich.console import Console
|
||||
from rich.logging import RichHandler
|
||||
from rich.traceback import install
|
||||
|
||||
from .env import env
|
||||
|
||||
console = Console(width=env.log.console_width, color_system="auto", force_terminal=True)
|
||||
|
||||
|
||||
def setup_logging() -> None:
|
||||
from aiogram.dispatcher import router # noqa: PLC0415
|
||||
|
||||
logging.basicConfig(
|
||||
level=env.log.level_external,
|
||||
format="",
|
||||
datefmt=None,
|
||||
handlers=[
|
||||
RichHandler(
|
||||
console=console,
|
||||
markup=True,
|
||||
rich_tracebacks=True,
|
||||
enable_link_path=False,
|
||||
tracebacks_show_locals=True,
|
||||
omit_repeated_times=False,
|
||||
show_time=env.log.show_time,
|
||||
tracebacks_suppress=[router],
|
||||
)
|
||||
],
|
||||
)
|
||||
install(console=console, show_locals=True)
|
||||
|
||||
|
||||
logger = logging.getLogger("healthbot")
|
||||
logger.setLevel(env.log.level)
|
||||
Reference in New Issue
Block a user