vidconf_db_up проверяется отдельным от основного пула соединением (NullPool, короткий таймаут) — иначе в момент исчерпания пула проверка сама встала бы в очередь и не отличила бы «БД лежит» от «пул занят». vidconf_db_pool_* читаются синхронно из engine.pool, без единого запроса к БД. metrics_endpoint больше не виснет и не падает при недоступном основном пуле: критичные gauge'и считаются первыми и не зависят от него, а vidconf_pipeline_sessions (по-прежнему через Depends(get_session) — тестовый харнесс подменяет её на savepoint-сессию) обёрнут таймаутом и try/except.
274 lines
15 KiB
Python
274 lines
15 KiB
Python
"""Метрики 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)
|