2026-07-01 23:15:35 +02:00
|
|
|
"""Monitoring: Prometheus metrics, structured JSON logging, health checks."""
|
|
|
|
|
|
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
|
|
|
|
from typing import Any
|
|
|
|
|
|
|
|
|
|
import structlog
|
|
|
|
|
from prometheus_client import (
|
|
|
|
|
CollectorRegistry,
|
|
|
|
|
Counter,
|
|
|
|
|
Gauge,
|
|
|
|
|
Histogram,
|
|
|
|
|
generate_latest,
|
|
|
|
|
)
|
|
|
|
|
from sqlalchemy import text
|
|
|
|
|
|
|
|
|
|
from app.core.db import get_engine
|
|
|
|
|
|
|
|
|
|
# ─── Prometheus Metrics Registry ───
|
|
|
|
|
|
|
|
|
|
REGISTRY = CollectorRegistry()
|
|
|
|
|
|
|
|
|
|
http_requests_total = Counter(
|
|
|
|
|
"leocrm_http_requests_total",
|
|
|
|
|
"Total HTTP requests by method, path, and status",
|
|
|
|
|
["method", "path", "status"],
|
|
|
|
|
registry=REGISTRY,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
http_request_duration_seconds = Histogram(
|
|
|
|
|
"leocrm_http_request_duration_seconds",
|
|
|
|
|
"HTTP request duration in seconds",
|
|
|
|
|
["method", "path"],
|
|
|
|
|
registry=REGISTRY,
|
|
|
|
|
buckets=(0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0),
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
db_pool_connections = Gauge(
|
|
|
|
|
"leocrm_db_pool_connections",
|
|
|
|
|
"Database connection pool size",
|
|
|
|
|
registry=REGISTRY,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
arq_jobs_total = Counter(
|
|
|
|
|
"leocrm_arq_jobs_total",
|
|
|
|
|
"Total ARQ background jobs by function and status",
|
|
|
|
|
["function", "status"],
|
|
|
|
|
registry=REGISTRY,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# ─── Structured Logging (structlog) ───
|
|
|
|
|
|
|
|
|
|
structlog.configure(
|
|
|
|
|
processors=[
|
|
|
|
|
structlog.contextvars.merge_contextvars,
|
|
|
|
|
structlog.processors.add_log_level,
|
|
|
|
|
structlog.processors.TimeStamper(fmt="iso"),
|
|
|
|
|
structlog.processors.JSONRenderer(),
|
|
|
|
|
],
|
|
|
|
|
wrapper_class=structlog.make_filtering_bound_logger(20), # INFO level
|
|
|
|
|
cache_logger_on_first_use=True,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
logger = structlog.get_logger("leocrm")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def get_logger(name: str = "leocrm") -> structlog.stdlib.BoundLogger:
|
|
|
|
|
"""Get a structured logger instance."""
|
|
|
|
|
return structlog.get_logger(name)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def record_request(
|
|
|
|
|
method: str,
|
|
|
|
|
path: str,
|
|
|
|
|
status_code: int,
|
|
|
|
|
duration_ms: float,
|
|
|
|
|
tenant_id: str | None = None,
|
|
|
|
|
) -> None:
|
|
|
|
|
"""Record an API request in metrics and structured log."""
|
|
|
|
|
http_requests_total.labels(method=method, path=path, status=str(status_code)).inc()
|
|
|
|
|
http_request_duration_seconds.labels(method=method, path=path).observe(
|
|
|
|
|
duration_ms / 1000.0
|
|
|
|
|
)
|
|
|
|
|
logger.info(
|
|
|
|
|
"api_request",
|
|
|
|
|
method=method,
|
|
|
|
|
path=path,
|
|
|
|
|
status=status_code,
|
|
|
|
|
duration_ms=round(duration_ms, 2),
|
|
|
|
|
tenant_id=tenant_id,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def record_error(
|
|
|
|
|
event: str,
|
|
|
|
|
method: str,
|
|
|
|
|
path: str,
|
|
|
|
|
status_code: int,
|
|
|
|
|
error: str,
|
|
|
|
|
traceback_str: str | None = None,
|
|
|
|
|
tenant_id: str | None = None,
|
|
|
|
|
) -> None:
|
|
|
|
|
"""Log an error with stacktrace and request context."""
|
|
|
|
|
logger.error(
|
|
|
|
|
event=event,
|
|
|
|
|
method=method,
|
|
|
|
|
path=path,
|
|
|
|
|
status=status_code,
|
|
|
|
|
error=error,
|
|
|
|
|
traceback=traceback_str,
|
|
|
|
|
tenant_id=tenant_id,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def update_db_pool_gauge() -> None:
|
|
|
|
|
"""Update the DB pool connections gauge from the current engine."""
|
|
|
|
|
engine = get_engine()
|
|
|
|
|
pool = engine.pool
|
|
|
|
|
db_pool_connections.set(pool.size() + pool.checkedout())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def record_arq_job(function: str, status: str) -> None:
|
|
|
|
|
"""Record an ARQ background job completion."""
|
|
|
|
|
arq_jobs_total.labels(function=function, status=status).inc()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def generate_metrics() -> bytes:
|
|
|
|
|
"""Generate Prometheus-formatted metrics output."""
|
|
|
|
|
update_db_pool_gauge()
|
|
|
|
|
return generate_latest(REGISTRY)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# ─── Health Checks ───
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def check_database() -> dict[str, Any]:
|
|
|
|
|
"""Check database connectivity (async, no blocking I/O)."""
|
|
|
|
|
try:
|
|
|
|
|
engine = get_engine()
|
|
|
|
|
async with engine.connect() as conn:
|
|
|
|
|
result = await conn.execute(text("SELECT 1"))
|
|
|
|
|
result.scalar()
|
|
|
|
|
return {"status": "up", "latency_ms": 0}
|
|
|
|
|
except Exception as e:
|
|
|
|
|
return {"status": "down", "error": str(e)}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def check_redis() -> dict[str, Any]:
|
|
|
|
|
"""Check Redis connectivity (async)."""
|
|
|
|
|
try:
|
|
|
|
|
import redis.asyncio as aioredis
|
|
|
|
|
|
|
|
|
|
from app.config import get_settings
|
|
|
|
|
|
|
|
|
|
settings = get_settings()
|
|
|
|
|
r = aioredis.from_url(settings.redis_url, decode_responses=True)
|
|
|
|
|
pong = await r.ping()
|
|
|
|
|
await r.aclose()
|
|
|
|
|
if pong:
|
|
|
|
|
return {"status": "up", "latency_ms": 0}
|
|
|
|
|
return {"status": "down", "error": "Redis returned False for PING"}
|
|
|
|
|
except Exception as e:
|
|
|
|
|
return {"status": "down", "error": str(e)}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def check_storage() -> dict[str, Any]:
|
|
|
|
|
"""Check storage directory accessibility."""
|
|
|
|
|
import os
|
|
|
|
|
|
|
|
|
|
from app.config import get_settings
|
|
|
|
|
|
|
|
|
|
settings = get_settings()
|
|
|
|
|
storage_path = getattr(settings, "storage_path", "/tmp")
|
|
|
|
|
try:
|
|
|
|
|
if os.path.exists(storage_path) and os.access(storage_path, os.W_OK): # noqa: ASYNC240
|
|
|
|
|
return {"status": "up", "path": storage_path}
|
|
|
|
|
return {"status": "down", "error": f"Storage path {storage_path} not writable"}
|
|
|
|
|
except Exception as e:
|
|
|
|
|
return {"status": "down", "error": str(e)}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def check_worker() -> dict[str, Any]:
|
|
|
|
|
"""Check if ARQ worker process is reachable (best-effort).
|
|
|
|
|
In test/dev mode this checks if Redis is available for the worker queue.
|
|
|
|
|
"""
|
|
|
|
|
try:
|
|
|
|
|
import redis.asyncio as aioredis
|
|
|
|
|
|
|
|
|
|
from app.config import get_settings
|
|
|
|
|
|
|
|
|
|
settings = get_settings()
|
|
|
|
|
r = aioredis.from_url(settings.redis_url, decode_responses=True)
|
|
|
|
|
# Check if arq queue key exists
|
2026-07-19 02:25:42 +02:00
|
|
|
queue_length = await r.zcard("arq:queue")
|
2026-07-01 23:15:35 +02:00
|
|
|
await r.aclose()
|
|
|
|
|
return {"status": "up", "queue_length": queue_length}
|
|
|
|
|
except Exception as e:
|
|
|
|
|
return {"status": "down", "error": str(e)}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def get_health_status() -> dict[str, Any]:
|
|
|
|
|
"""Run all health checks and return aggregated status."""
|
|
|
|
|
db_result = await check_database()
|
|
|
|
|
redis_result = await check_redis()
|
|
|
|
|
storage_result = await check_storage()
|
|
|
|
|
worker_result = await check_worker()
|
|
|
|
|
|
|
|
|
|
all_up = all(
|
|
|
|
|
c["status"] == "up"
|
|
|
|
|
for c in [db_result, redis_result, storage_result, worker_result]
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
if all_up:
|
|
|
|
|
overall_status = "healthy"
|
|
|
|
|
elif db_result["status"] == "down":
|
|
|
|
|
overall_status = "degraded"
|
|
|
|
|
else:
|
|
|
|
|
overall_status = "degraded"
|
|
|
|
|
|
|
|
|
|
return {
|
|
|
|
|
"status": overall_status,
|
|
|
|
|
"version": "1.0.0",
|
|
|
|
|
"checks": {
|
|
|
|
|
"database": db_result,
|
|
|
|
|
"redis": redis_result,
|
|
|
|
|
"storage": storage_result,
|
|
|
|
|
"worker": worker_result,
|
|
|
|
|
},
|
|
|
|
|
}
|