Newer
Older
hard-panel / panel / backend / app / mcp.py
"""MCP-сервер панели: любые ИИ-агенты читают состояние серверов по /mcp.

Транспорт — streamable HTTP (stateless + JSON-ответы), приложение
монтируется в FastAPI в app.main на /mcp. Клиент Claude Code:

    claude mcp add --transport http ghard https://panel.example.com/mcp

Инструкции агенту (AGENT_INSTRUCTIONS) отдаются и как instructions
инициализации, и как prompt agent_guide — паттерн reference-реализации
gnexus-creds (data_api). Тулы — чтение плюс управление серверами и
сервисами (запись по явной просьбе человека).
"""

import json
import time
from datetime import datetime, timedelta, timezone

from fastapi import HTTPException, Request
from pydantic import ValidationError
from starlette.responses import Response
from mcp.server.fastmcp import FastMCP

from app.api.events import list_events as _list_events
from app.api.servers import (
    create_server as _create_server,
    delete_server as _delete_server,
    get_server,
    list_servers,
)
from app.api.services import (
    ServiceCreate as _ServiceCreate,
    create_service as _create_service,
    delete_service as _delete_service,
    list_services as _list_services,
    service_incidents as _service_incidents,
)
from app.api.shares import list_shares as _list_shares
from app.db import get_db
from app.models import ServerCreate as _ServerCreate

AGENT_INSTRUCTIONS = """\
# GHard Monitor — инструкция для ИИ-агента

## Назначение
Этот сервер — MCP-интерфейс панели мониторинга GHard Monitor
(Gnexus Hardware Monitor). Основная твоя роль — наблюдатель-диагност:
отвечаешь пользователю о состоянии его серверов и сервисов. Читающие
тулы дают полную картину поверх наблюдения; записывающие тулы
(создание/удаление серверов и сервисов) — только по явной просьбе
человека: панель общая, её изменения видны всем.
Перезапускать агентов и сервисы ты не можешь — это за пределами панели.

## Данные
- **Серверы** — машины с установленным агентом hard-monitor: CPU/RAM/swap,
  load, диски по маунтам, сеть (счёт с агента), процессы, docker-контейнеры,
  апптайм, заметка пользователя.
- **Сервисы** — health-чеки внешних HTTP-эндпоинтов (gnexus-сервисы:
  /health и похожие). Панель пробует их по расписанию и хранит историю
  откликов + журнал моментов недоступности.
- **Сетевые хранилища** — примонтированные NAS-шары и их заполнение.
- **Журнал событий** — всё, что панель зафиксировала: offline/пороги/
  контейнеры/сервисы (хранение 90 дней).

## Язык
Отвечай пользователю на его языке. Вывод тулов технический (статусы
и события — английские слаги), переводи их в человеческую речь сам.

## Карта инструментов — что вызывать
| Вопрос пользователя | Инструмент |
|---|---|
| «как мои серверы?» / «что вообще происходит?» | `panel_overview` — одним вызовом: все серверы + сервисы + хранилища |
| «что с web-01?» / «что жрёт память?» / «почему тормозит?» | `server_details` (id — из `panel_overview`) |
| «растёт ли диск на db-01 за сутки?» / тренд | `server_history` |
| «сервис падал?», «с какими сервисами проблемы?» | `services_overview`, потом `service_incidents` (id — из overview) |
| «что происходило ночью / за последние часы?» | `events_recent` |
| «добавь сервер/сервис», «удали web-old» | тулы из «Управление (запись)» |

Порядок для составных вопросов: обзор → выбранная детализация. `panel_overview`
дешёвый — начинай с него, если не уверен.

## Управление (запись)
Четыре тула меняют каталог. Запись лимитирована: ≤20 записывающих вызовов
за 10 минут, сверх того возвращается `{"error": 429}`.
- `server_create(name, hostname="", interval=30)` — регистрирует сервер и
  возвращает ключ агента `ghm_…`. Ключ показывается **ровно один раз —
  сохрани его сразу** и передай пользователю (в .env агента). До установки
  агента сервер остаётся `pending`.
- `server_delete(server_id)` — удаляет сервер **необратимо**: его метрики
  и журнал событий стираются каскадом.
- `service_create(name, url)` — добавляет health-чек; первая проба
  выполняется сразу, статус читай через `services_overview`.
- `service_delete(service_id)` — удаляет сервис; история его проб стирается.

Правила записи:
- удалять — **только по явной просьбе пользователя** («удали», «убери»),
  не по собственным соображениям; перед удалением назови, что сотрётся;
- ошибки записи приходят как данные: `{"error": 404|422|429, "detail": "…"}`
  — прочти detail и действуй по нему;
- результат создания (id, ключ, URL) возвращай пользователю — это его данные.

## Словарь статусов
- Сервер: `online` (пакеты идут) / `offline` (панель не получает метрики —
  агент умер, сеть или машина выключена) / данных нет.
- Сервис: `up` (здоров) / `degraded` (работает, но сообщает о деградации) /
  `down` (не отвечает) / `pending` (данных ещё нет). Поле `mode`
  (maintenance / readonly / draining) — это **режим работы**, а не падение;
  не пугай пользователя режимом.
- Журнал: severity `info` / `warning` / `critical`.

## Словарь событий (поле type в журнале)
- `server_offline` / `server_online` — сервер перестал/возобновил слать метрики.
- `cpu_high` / `ram_high` / `swap_high` / `disk_high` / `load_high` — ресурс
  превысил порог (порог и значение — в payload); `*_recovered` — вернулся
  в норму (гистерезис: открытие и закрытие — разные пороги).
- `container_added` / `container_removed` / `container_started` /
  `container_exited` — docker-контейнер появился/исчез/запустился/завершился
  (в payload: name, image, exit_code если есть).
- `service_down` / `service_degraded` / `service_recovered` — сервис
  перестал отвечать / деградировал / восстановился.
- Пороги сжимаются в человеческую речь, например:
  `cpu 95% (≥90%)`, `disk / 98% (≥90%)`, `load1 2.1×cores (≥2.0×)`.

## Правила
- ID серверов и сервисов бери из обзорных тулов — угадывание запрещено.
- `server_history(hours)`: разумно 1 (детали), 6, 24 (сутки) или 168 (неделя;
  глубже недели данных нет — история метрик хранится 7 дней).
- Числа передавай как есть из тулов (проценты, МБ/с, ms) — не пересчитывай
  единицы самостоятельно.
- Если сервер offline — не паникуй и не повторяй вопросы тулами; скажи
  последнее известное состояние (в `server_details`) и время последнего пакета.
- Пользователь просит «что-нибудь сделать» (перезапустить, поправить) —
  перезапуски за пределами панели, а создание и удаление — тулы из
  «Управление (запись)», только по явной просьбе.

## Подключение (для администратора панели)
```
claude mcp add --transport http ghard https://panel.example.com/mcp
```
Transport — streamable HTTP. Авторизация — Bearer-ключ в заголовке
`Authorization`:
- персональный токен из страницы «MCP-ключи» веб-интерфейса панели
  (`Authorization: Bearer mcp_…`) — ключ принадлежит вашему пользователю,
  отзыв виден на той же странице;
- `GHARD_ADMIN_TOKEN` — супер-токен администратора, действует как раньше.
"""


mcp = FastMCP(
    "GHard Monitor",
    instructions=AGENT_INSTRUCTIONS,
    stateless_http=True,  # без сессий — работает за reverse proxy
    json_response=True,
    streamable_http_path="/",
)


@mcp.prompt()
def agent_guide() -> str:
    """📘 Гайд: как отвечать на вопросы о серверах и сервисах через GHard Monitor."""
    return AGENT_INSTRUCTIONS


def _dumps(data) -> str:
    return json.dumps(data, ensure_ascii=False, indent=1)


def _err(code: int, detail: str) -> str:
    """Ошибка записывающего тула как данные (канон handbook mcp.md)."""
    return json.dumps({"error": code, "detail": detail}, ensure_ascii=False)


# Rate-limit записи (канон mcp.md: чувствительные операции). Окно в памяти
# процесса, глобальное — тула не видит личность звонящего (FastMCP этой
# версии не пробрасывает HTTP-контекст), для защиты от спама хватает.
_WRITE_WINDOW = 600.0
_WRITE_MAX = 20
_write_log: list[float] = []


def _write_gate() -> str | None:
    """None — проход; строка — JSON-ошибка 429 (лимит записи израсходован)."""
    now = time.monotonic()
    while _write_log and now - _write_log[0] > _WRITE_WINDOW:
        del _write_log[0]
    if len(_write_log) >= _WRITE_MAX:
        return _err(
            429,
            f"too many write calls (>={_WRITE_MAX} за {int(_WRITE_WINDOW / 60)} мин) — повторите позже",
        )
    _write_log.append(now)
    return None


def _worst_disk(disks: list[dict] | None) -> dict | None:
    if not disks:
        return None
    return max(disks, key=lambda d: d.get("percent", 0))


@mcp.tool()
async def panel_overview() -> str:
    """Обзор всего: серверы (статус, CPU/RAM/диски, сеть, последний пакет),
    сервисы (health), сетевые хранилища. Самый быстрый способ ответить на
    «как мои серверы?» / «что вообще происходит?»."""
    servers = await list_servers()
    online = sum(1 for s in servers if s["status"] == "online")
    lines = [f"Серверов: {len(servers)}, онлайн: {online}, не в сети: {len(servers) - online}"]
    for s in servers:
        m = s.get("metrics")
        if m:
            worst = _worst_disk(m.get("disks"))
            disk = f", диск {worst['mount']} {worst['percent']:.0f}%" if worst else ""
            lines.append(
                f"#{s['id']} {s['name']} [{s['status']}] CPU {m['cpu']:.0f}%, "
                f"RAM {m['ram']['used'] / 2**30:.1f}/{m['ram']['total'] / 2**30:.1f} GiB{disk}, "
                f"сеть {m['net_in_mbs']:.2f}/{m['net_out_mbs']:.2f} MB/s, "
                f"пакет {s['last_seen']}"
            )
        else:
            status = s['status']
            lines.append(f"#{s['id']} {s['name']} [{status}] — данных ещё нет")
    services = await _list_services()
    if services:
        bad = sum(1 for x in services if x["status"] in ("down", "degraded"))
        lines.append(f"Сервисов: {len(services)}, проблемных (down/degraded): {bad}")
        for svc in services:
            report = svc.get("report") or {}
            identity = ""
            if report.get("version"):
                identity = f", {report.get('service') or '?'} {report['version']}"
            mode = f", режим {report['mode']}" if report.get("mode") else ""
            msg = f" ({svc['message']})" if svc.get("message") else ""
            lines.append(
                f"#{svc['id']} {svc['name']} [{svc['status']}] {svc['latency_ms'] or '—'} ms{identity}{mode}{msg}"
            )
    shares = await _list_shares()
    if shares:
        lines.append(f"Сетевые хранилища: {len(shares)}")
        for sh in shares:
            if sh["total"]:
                percent = sh["used"] / sh["total"] * 100
                lines.append(
                    f"#{sh['id']} {sh['name']} [{sh['status']}] "
                    f"{sh['used'] / 2**30:.1f}/{sh['total'] / 2**30:.1f} GiB ({percent:.0f}%), {sh['path']}"
                )
            else:
                lines.append(f"#{sh['id']} {sh['name']} [{sh['status']}] — нет данных, {sh['path']}")
    return "\n".join(lines)


@mcp.tool()
async def server_details(server_id: int) -> str:
    """Полное текущее состояние одного сервера: метрики, диски, топ процессов,
    docker-контейнеры, заметка. id бери из panel_overview."""
    try:
        server = await get_server(server_id)
    except Exception:
        return f"Сервер #{server_id} не найден"
    return _dumps(server)


@mcp.tool()
async def server_history(server_id: int, hours: float = 1.0, max_points: int = 200) -> str:
    """История метрик сервера за `hours` часов, прореженная до `max_points`.

    Возвращает точки: ts, cpu, ram%, swap, load, сеть MB/s, диски.
    История метрик хранится 7 дней — глубже точки удаляются.
    """
    since = (datetime.now(timezone.utc) - timedelta(hours=hours)).isoformat()
    db = get_db()
    cursor = await db.execute(
        """SELECT ts, cpu, load1, load5, load15, ram_used, ram_total,
                  swap_used, swap_total, uptime, net_in_mbs, net_out_mbs, disks_json
           FROM metrics WHERE server_id = ? AND ts >= ?
           ORDER BY ts DESC LIMIT 5000""",
        (server_id, since),
    )
    rows = await cursor.fetchall()
    points = [
        {
            "ts": row["ts"],
            "cpu": row["cpu"],
            "load": [row["load1"], row["load5"], row["load15"]],
            "ram_percent": round(row["ram_used"] / row["ram_total"] * 100, 1)
            if row["ram_total"]
            else None,
            "swap_percent": round(row["swap_used"] / row["swap_total"] * 100, 1)
            if row["swap_total"]
            else None,
            "uptime": row["uptime"],
            "net_in_mbs": row["net_in_mbs"],
            "net_out_mbs": row["net_out_mbs"],
            "disks": json.loads(row["disks_json"]),
        }
        for row in reversed(rows)
    ]
    if len(points) > max_points:
        step = len(points) / max_points
        points = [points[int(i * step)] for i in range(max_points)]
    return _dumps({"server_id": server_id, "hours": hours, "points": points})


# --- Сервисы (health-чеки) -----------------------------------------------------


@mcp.tool()
async def services_overview() -> str:
    """Все наблюдаемые сервисы (health-чеки): state, latency, HTTP-код,
    self-report (версия/окружение/режим). id нужен для service_incidents."""
    services = await _list_services()
    lines = [f"Сервисов: {len(services)}"]
    for svc in services:
        report = svc.get("report") or {}
        identity = []
        for key in ("service", "version", "environment"):
            if report.get(key):
                identity.append(str(report[key]))
        id_str = f" — {' '.join(identity)}" if identity else ""
        mode = f", режим {report['mode']}" if report.get("mode") else ""
        code = f" код {svc['code']}" if svc.get("code") is not None else ""
        msg = f" ({svc['message']})" if svc.get("message") else ""
        lines.append(
            f"#{svc['id']} {svc['name']} [{svc['status']}] "
            f"{svc['latency_ms'] or '—'} ms{code}{mode}{msg}{id_str}"
        )
    return "\n".join(lines)


@mcp.tool()
async def service_incidents(service_id: int, limit: int = 20) -> str:
    """Журнал моментов недоступности сервиса: периоды down+degraded.

    Возвращает: start, end (null = идёт сейчас), duration_s,
    worst (худшее состояние в периоде), code, message (причина из первой точки).
    """
    try:
        incidents = await _service_incidents(service_id, limit=limit)
    except Exception:
        return f"Сервис #{service_id} не найден"
    return _dumps(incidents)


# --- Журнал событий ------------------------------------------------------------


@mcp.tool()
async def events_recent(limit: int = 30, severity: str = "") -> str:
    """Хвост единого журнала событий (90 дней хранения).

    severity — необязательный фильтр: info | warning | critical.
    Возвращает: id, ts, type, severity, message, data (payload события).
    """
    if limit < 1 or limit > 1000:
        return "limit: 1..1000"
    events = await _list_events(limit=min(limit, 1000))
    if severity:
        events = [e for e in events if e.get("severity") == severity]
    return _dumps(events)


# --- Управление (запись): серверы и сервисы ------------------------------------


@mcp.tool()
async def server_create(name: str, hostname: str = "", interval: int = 30) -> str:
    """Зарегистрировать сервер и получить ключ агента. Ключ (ghm_…) показан
    ОДИН раз — сразу сохрани его для .env агента. До установки агента сервер
    в статусе pending."""
    if gate := _write_gate():
        return gate
    if not name.strip():
        return _err(422, "name обязателен")
    if not 5 <= interval <= 3600:
        return _err(422, "interval: 5..3600 секунд")
    try:
        created = await _create_server(
            _ServerCreate(name=name.strip(), hostname=hostname.strip(), interval=interval)
        )
    except ValidationError as exc:
        return _err(422, str(exc))
    return (
        f"Сервер #{created['id']} «{created['name']}» создан.\n"
        f"Ключ агента (показан один раз — сохрани немедленно):\n"
        f"{created['key']}\n"
        f"До установки агента (PANEL_URL + SERVER_KEY + INTERVAL в .env) сервер в статусе pending."
    )


@mcp.tool()
async def server_delete(server_id: int) -> str:
    """Удалить сервер НЕОБРАТИМО: его метрики и журнал событий стираются
    каскадом. Вызывай только по явной просьбе пользователя."""
    if gate := _write_gate():
        return gate
    try:
        await _delete_server(server_id)
    except HTTPException as exc:
        return _err(exc.status_code, str(exc.detail))
    return f"Сервер #{server_id} удалён; его метрики и журнал стёрты безвозвратно."


@mcp.tool()
async def service_create(name: str, url: str) -> str:
    """Добавить health-чек сервиса (name до 100 симв., url до 500).
    Первая проба выполняется сразу; статус — через services_overview."""
    if gate := _write_gate():
        return gate
    if not url.strip():
        return _err(422, "url обязателен")
    try:
        created = await _create_service(_ServiceCreate(name=name.strip(), url=url.strip()))
    except ValidationError:
        return _err(422, "name: 1..100 символов, url: 1..500 символов")
    return (
        f"Сервис #{created['id']} «{created['name']}» добавлен: {created['url']}\n"
        f"Первая проба выполнена (или доберёт фоновый цикл) — статус читай в services_overview."
    )


@mcp.tool()
async def service_delete(service_id: int) -> str:
    """Удалить сервис НЕОБРАТИМО: история его проб стирается. Вызывай только
    по явной просьбе пользователя."""
    if gate := _write_gate():
        return gate
    try:
        await _delete_service(service_id)
    except HTTPException as exc:
        return _err(exc.status_code, str(exc.detail))
    return f"Сервис #{service_id} удалён; история его проб стёрта безвозвратно."


mcp_app = mcp.streamable_http_app()

# ASGI-эндпоинт MCP: регистрируется прямо в FastAPI на /mcp (см. app.main),
# потому что Mount() не совпадает с путём без хвостового слэша, а lifespan
# смонтированного sub-app не запускается сам — менеджер сессий стартует
# в lifespan app.main (mcp.session_manager.run()).
mcp_asgi = mcp_app.router.routes[0].app


async def mcp_endpoint(request: Request) -> "Response":
    """Мост FastAPI → ASGI streamable-http-app.

    FastAPI валидирует сигнатуру endpoint'а: голый ASGI-объект (scope,
    receive, send) был бы прочитан как query-параметры → 422. Поэтому мост:
    забираем тело, вручную докручиваем ASGI-диалог, ответ собираем из
    сообщений send (stateless + json_response=True — ответ всегда один
    JSON, SSE-стримов тут нет).
    """
    body = await request.body()
    received = False
    start: dict = {}
    chunks: list[bytes] = []

    async def receive():
        nonlocal received
        if received:
            return {"type": "http.disconnect"}
        received = True
        return {"type": "http.request", "body": body, "more_body": False}

    async def send(message) -> None:
        if message["type"] == "http.response.start":
            start.update(message)
        else:
            chunks.append(message.get("body", b""))

    await mcp_asgi(request.scope, receive, send)
    response = Response(
        content=b"".join(chunks) if chunks else b"", status_code=start["status"]
    )
    response.raw_headers.extend(start.get("headers") or [])
    return response