"""Метрики 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` источник и вовсе синхронный — настройки уже в памяти процесса). 🔴 Метрики о состоянии основного пула БД (`vidconf_db_up`, `vidconf_db_pool_*`) обязаны читаться БЕЗ обращения к самому пулу — иначе в момент его исчерпания (см. `.forcc/session-results/32-loadtest-07-08-debug.md`) эндпоинт метрик падал бы вместе со всем остальным ровно тогда, когда нужнее всего. `vidconf_db_pool_*` — синхронный снимок `engine.pool` (см. `core/db.py::db_pool_stats`), `vidconf_db_up` — отдельное соединение вне основного пула (`core/db.py::check_db_up`). `_refresh_pipeline_sessions_gauge` по-прежнему ходит через основной пул (`Depends(get_session)`, тестовый харнесс подменяет её на savepoint-сессию — см. `tests/conftest.py`; развести полностью, как `vidconf_db_up`, значило бы переделывать харнесс ради того же эффекта — цена не оправдана, см. прецедент `f7c4fb4`/session 32), но обёрнута таймаутом и try/except, чтобы её недоступность не роняла остальные метрики. """ import asyncio 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 check_db_up, db_pool_checked_out, 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",), ) # Сколько ждать основной пул под этой конкретной метрикой, прежде чем # сдаться и оставить прежнее значение gauge. Меньше `db_pool_timeout` (10с, # `core/config.py`) — Prometheus скрейпит раз в 15с, и эта метрика не должна # в одиночку съедать бюджет всего окна scrape. _PIPELINE_GAUGE_TIMEOUT_S = 2.0 async def _refresh_pipeline_sessions_gauge(session: AsyncSession) -> None: """Пересчитать `vidconf_pipeline_sessions` по всем статусам `pipeline_status`. Ходит через основной пул (`session` — из `Depends(get_session)`, см. докстринг модуля про ограничения тестового харнесса). Если пул занят или БД недоступна, запрос не должен держать весь `/metrics` — таймаут короче `db_pool_timeout`, ошибка гасится, gauge остаётся на прежнем значении (не обнуляется — обнулять его при недоступности БД так же неверно, как считать сеансы пропавшими). """ try: counts = await asyncio.wait_for( ConferenceSessionRepository(session).count_by_pipeline_status(), timeout=_PIPELINE_GAUGE_TIMEOUT_S, ) except Exception: # noqa: BLE001 return 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) # --- Доступность БД и занятость основного пула (сессия 33) ----------------- # # Ранний сигнал важнее самого факта отказа: в инциденте 07.08 пул заполнялся # постепенно (`idle in transaction` 3→8→16→26→35→39→40 участников) — # `vidconf_db_pool_checked_out` показал бы это задолго до первого 500. # Обе метрики читаются без обращения к основному пулу (см. докстринг модуля # и `core/db.py`), поэтому доступны и в момент, когда сам пул исчерпан. DB_UP = Gauge( "vidconf_db_up", "Доступность БД (1/0) — проверяется отдельным соединением вне основного пула", ) DB_POOL_SIZE = Gauge( "vidconf_db_pool_size", "Настроенный размер основного пула БД без overflow (db_pool_size)", ) DB_POOL_MAX_OVERFLOW = Gauge( "vidconf_db_pool_max_overflow", "Настроенный максимум overflow-соединений сверх db_pool_size (db_max_overflow)", ) DB_POOL_CHECKED_OUT = Gauge( "vidconf_db_pool_checked_out", "Число соединений основного пула БД, занятых прямо сейчас (в пуле + overflow)", ) async def _refresh_db_up_gauge() -> None: """Пересчитать `vidconf_db_up` отдельным от основного пула соединением.""" DB_UP.set(1 if await check_db_up() else 0) def _refresh_db_pool_gauges() -> None: """Пересчитать gauge'и занятости основного пула — синхронно, без I/O.""" settings = get_settings() DB_POOL_SIZE.set(settings.db_pool_size) DB_POOL_MAX_OVERFLOW.set(settings.db_max_overflow) DB_POOL_CHECKED_OUT.set(db_pool_checked_out()) # --- Занятость пула Redis (сессия 33, второй потолок из session 32) -------- # # Тот же класс отказа, что и у пула БД: каждое WS-подключение комнаты держит # pub/sub-соединение всё время, пока участник в конференции (`core/redis.py`, # `redis_max_connections`). Снимок — синхронный (атрибуты пула в памяти # процесса redis-py), Redis для этого спрашивать не нужно. REDIS_POOL_IN_USE = Gauge( "vidconf_redis_pool_in_use", "Число занятых соединений пула Redis прямо сейчас", ) REDIS_POOL_MAX = Gauge( "vidconf_redis_pool_max_connections", "Настроенный максимум соединений пула Redis (redis_max_connections)", ) def _refresh_redis_pool_gauges() -> None: """Пересчитать gauge'и занятости пула Redis — синхронно, без I/O.""" pool = redis_client.connection_pool REDIS_POOL_IN_USE.set(len(pool._in_use_connections)) # noqa: SLF001 REDIS_POOL_MAX.set(pool.max_connections) # --- 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с), нагрузка пренебрежимо мала. Порядок важен: метрики о состоянии основного пула БД (`_refresh_db_up_gauge`, `_refresh_db_pool_gauges`) считаются первыми и не зависят от самого пула (см. докстринг модуля) — они гарантированно попадут в ответ, даже если следующий за ними `_refresh_pipeline_sessions_gauge` (основной пул) зависнет или упадёт под нагрузкой. """ await _refresh_db_up_gauge() _refresh_db_pool_gauges() _refresh_redis_pool_gauges() 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)