"""Метрики Prometheus: латентность HTTP + gauge'и пайплайна, очередей и железа. `GET /metrics` — без авторизации (снаружи закрывается на уровне nginx, вне периметра backend, см. `docs/deploy/scaling.md`/monitoring-часть devops): Prometheus-серверы традиционно ходят напрямую в контейнер по внутренней сети, а не через публичный `/api/`-гейтвей. Gauge'и `vidconf_pipeline_sessions`/`vidconf_celery_queue_depth`/ `vidconf_host_info` намеренно НЕ обновляются фоновой задачей — значения пересчитываются прямо в обработчике запроса при каждом scrape (см. докстринг `metrics_endpoint`), поэтому их асинхронные источники (БД, Redis) можно опросить обычным `await` вместо реализации синхронного `prometheus_client.registry.Collector` (у `vidconf_host_info` источник и вовсе синхронный — настройки уже в памяти процесса). """ import time from collections.abc import Awaitable, Callable from fastapi import APIRouter, Depends, Request, Response from prometheus_client import CONTENT_TYPE_LATEST, Gauge, Histogram, generate_latest from sqlalchemy.ext.asyncio import AsyncSession from starlette.routing import Match from core.config import get_settings from core.db import get_session from core.redis import redis_client from models.session import PIPELINE_STATUSES from repositories.conferences import ConferenceSessionRepository router = APIRouter() # --- Латентность HTTP-запросов по маршрутам -------------------------------- HTTP_REQUEST_DURATION_SECONDS = Histogram( "vidconf_http_request_duration_seconds", "Латентность HTTP-запросов backend по маршрутам", labelnames=("method", "path", "status"), ) async def prometheus_latency_middleware( request: Request, call_next: Callable[[Request], Awaitable[Response]] ) -> Response: """Замерить латентность запроса и записать в `HTTP_REQUEST_DURATION_SECONDS`. Метка `path` — шаблон маршрута (`/api/v1/conferences/{conference_id}`), а не сырой URL: иначе каждый UUID/slug в пути породил бы собственную серию меток (неограниченная кардинальность). Шаблон резолвится постфактум поиском совпавшего маршрута среди `request.app.routes` (FastAPI/Starlette не кладёт его в `request.scope` до входа в сам эндпоинт, а `call_next` оборачивает вызов целиком) — тот же приём, что использует `starlette.routing.Router` внутри себя для диспетчеризации. """ start = time.perf_counter() response = await call_next(request) duration = time.perf_counter() - start path_template = _match_route_path(request) HTTP_REQUEST_DURATION_SECONDS.labels( method=request.method, path=path_template, status=str(response.status_code) ).observe(duration) return response def _match_route_path(request: Request) -> str: """Найти шаблон пути совпавшего маршрута; сырой `request.url.path`, если не найден (404).""" for route in request.app.routes: match, _ = route.matches(request.scope) if match == Match.FULL: return getattr(route, "path", request.url.path) return request.url.path # --- Gauge числа сеансов по статусу пайплайна ------------------------------- PIPELINE_SESSIONS = Gauge( "vidconf_pipeline_sessions", "Число сеансов конференций в каждом статусе пайплайна пост-обработки", labelnames=("status",), ) async def _refresh_pipeline_sessions_gauge(session: AsyncSession) -> None: """Пересчитать `vidconf_pipeline_sessions` по всем статусам `pipeline_status`.""" counts = await ConferenceSessionRepository(session).count_by_pipeline_status() for status in PIPELINE_STATUSES: PIPELINE_SESSIONS.labels(status=status).set(counts.get(status, 0)) # --- Gauge глубины очередей Celery (Redis) ---------------------------------- CELERY_QUEUES = ("transcription", "summarize", "notify", "celery") """Очереди, за которыми следим (`workers/celery_app.py::app.conf.task_routes`, `docs/deploy/scaling.md`): выделенные `transcription`/`summarize`/`notify` + дефолтная `celery` (обслуживающие задачи без явного маршрута).""" CELERY_QUEUE_DEPTH = Gauge( "vidconf_celery_queue_depth", "Число задач, ожидающих обработки в очереди Celery (redis LLEN)", labelnames=("queue",), ) async def _refresh_celery_queue_depth_gauge() -> None: """Пересчитать `vidconf_celery_queue_depth` по всем отслеживаемым очередям. Список Redis, лежащий за очередью Celery, называется так же, как сама очередь (транспорт `kombu` с брокером `redis` кладёт задачи в список по имени очереди) — `LLEN` даёт точную глубину backlog'а на момент scrape. """ for queue in CELERY_QUEUES: depth = await redis_client.llen(queue) CELERY_QUEUE_DEPTH.labels(queue=queue).set(depth) # --- Info-метрика обнаруженного железа (install.sh, ADR-004) --------------- HOST_INFO = Gauge( "vidconf_host_info", "Обнаруженное установщиком железо (info-метрика, значение всегда 1, данные в лейблах)", labelnames=("cpus", "ram_mb", "gpu_name", "vram_mb"), ) def _hw_label(value: int | None) -> str: """`None` (install.sh не запускался либо GPU не обнаружен) → «—», не пустая строка. Пустой лейбл Grafana отрисовала бы как пустой текст на панели — прочерк однозначно читается как «не определено». """ return str(value) if value is not None else "—" def _refresh_host_info_gauge() -> None: """Пересчитать `vidconf_host_info` по текущим `HW_*` настройкам (`core/config.py`). Источник — `install.sh`, который пишет `HW_CPUS`/`HW_RAM_MB`/ `HW_GPU_NAME`/`HW_VRAM_MB` в `.env` при установке (ADR-004). Живые CPU/RAM/диск хоста в дашборде «Хост и контейнеры» берутся из node-exporter напрямую — здесь только то, чего node-exporter не знает (GPU), плюс дублирование CPU/RAM install.sh для сверки. """ settings = get_settings() HOST_INFO.labels( cpus=_hw_label(settings.hw_cpus), ram_mb=_hw_label(settings.hw_ram_mb), gpu_name=settings.hw_gpu_name or "—", vram_mb=_hw_label(settings.hw_vram_mb), ).set(1) @router.get("/metrics") async def metrics_endpoint(session: AsyncSession = Depends(get_session)) -> Response: """Отдать метрики Prometheus в формате text exposition. Gauge'и пересчитываются прямо здесь (а не по расписанию/периодическим коллектором) — значение в ответе всегда актуально на момент scrape, ценой одного SELECT (группировка по `pipeline_status`) и `LLEN` на каждую из 4 отслеживаемых очередей per запрос — Prometheus скрейпит редко (обычно раз в 15–30с), нагрузка пренебрежимо мала. """ await _refresh_pipeline_sessions_gauge(session) await _refresh_celery_queue_depth_gauge() _refresh_host_info_gauge() return Response(content=generate_latest(), media_type=CONTENT_TYPE_LATEST)