# ============================================ # Projet : API-KA # Fichier : src/api/routes/stats.py # Node : m3u96b # Author : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Date : 2026-08-19 # ============================================ """Routes /api/stats : tableau de bord analytique de la plateforme + rapports PDF. Contrat commun Groupe KA v2 (src/api/web/ka/stats/SPEC.md) : - ``GET /api/stats/dashboard?period=…`` → JSON (KPI + sparklines, jauges, séries, multi-séries, empilées, répartitions avec deltas, distributions, heatmap calendrier + horaire 7×24, tableaux, records) calculé depuis les données réelles : table ``api_requests`` (journal du middleware), ``collection_runs`` et ``connector_health``. - ``GET /api/stats/report?period=…&mode=complet|synthese|tendances| repartitions|donnees`` → PDF Groupe-KA (moteur kapdf v2, 5 rapports — un mode inconnu retombe sur ``complet``). - ``GET /api/stats/ecosystem-report?period=…&mode=…`` → UN SEUL PDF consolidant les 12 plateformes de l'écosystème : les dashboards des plateformes sœurs sont lus en parallèle (httpx, 15 s), le sien en local ; une plateforme injoignable devient une section « données indisponibles ». ``mode`` optionnel (mêmes 5 modes, inconnu ou absent → complet). Cache serveur : 5 minutes par période. Aucune stat inventée : une section sans données est simplement absente (le front affiche « Pas encore mesuré »). """ from __future__ import annotations import asyncio import datetime import json import threading import time from pathlib import Path from typing import Any from zoneinfo import ZoneInfo import httpx from fastapi import APIRouter, Body, Depends, HTTPException, Query, Response from sqlalchemy import func, select from sqlalchemy.orm import Session from src.api import ecopdf, kapdf from src.api.routes import envelope from src.database.db import get_db from src.database.models import ApiRequest, CollectionRun, ConnectorHealth router = APIRouter(prefix="/api/stats", tags=["stats"]) TZ = ZoneInfo("America/Toronto") CACHE_TTL = 300.0 # ≥ 5 min (SPEC) PLATFORM_ID = "api-ka" SITE = { "wordmark": "API·Ka", "accent": "#3b5bdb", "domain": "www.api-ka.com", "tagline": "Plateforme centrale de l'écosystème Groupe KA", } SERVICE_NAMES = { "louka": "Lou·Ka", "immoka": "Immo·Ka", "foodka": "Food·Ka", "autoka": "Auto·Ka", "fabrika": "Fabri·Ka", "restoka": "Resto·Ka", "sortika": "Sorti·Ka", "creaka": "Créa·Ka", "jobka": "Job·Ka", } PERIOD_LABELS = { "auj": "aujourd'hui", "7j": "7 jours", "30j": "30 jours", "3m": "3 mois", "6m": "6 mois", "12m": "12 mois", "annee": "année en cours", "tout": "toute la période", } PERIOD_DAYS = {"7j": 7, "30j": 30, "3m": 91, "6m": 182, "12m": 365} RUN_STATUS_FR = {"success": "succès", "failed": "échec", "retried": "relancé"} HEALTH_FR = {"ok": "OK", "degraded": "dégradé", "broken": "en panne", "stale": "périmé"} LAT_BINS = [ (0, 25, "< 25 ms"), (25, 50, "25-50 ms"), (50, 100, "50-100 ms"), (100, 200, "100-200 ms"), (200, 500, "200-500 ms"), (500, 1000, "0,5-1 s"), (1000, 2000, "1-2 s"), (2000, float("inf"), "> 2 s"), ] DUR_BINS = [ (0, 1, "< 1 s"), (1, 5, "1-5 s"), (5, 15, "5-15 s"), (15, 30, "15-30 s"), (30, 60, "30-60 s"), (60, float("inf"), "> 60 s"), ] _cache: dict[tuple, tuple[float, dict]] = {} _cache_lock = threading.Lock() # ------------------------------------------------------ rapport écosystème ECO_CACHE_TTL = 600.0 # 10 min par (period, mode) ECO_FETCH_TIMEOUT = 15.0 _eco_cache: dict[tuple[str, str], tuple[float, bytes, str]] = {} _eco_lock = threading.Lock() _ECOSYSTEM_JSON = Path(__file__).resolve().parents[1] / "web" / "ka" / "ecosystem.json" # Les 11 plateformes de données (le hub groupe-ka n'expose pas de dashboard). ECO_SITES: list[dict] = [ s for s in json.loads(_ECOSYSTEM_JSON.read_text(encoding="utf-8"))["sites"] if s["id"] != "groupe-ka" ] # ---------------------------------------------------------------- périodes def _resolve_period( period: str, date_from: datetime.date | None, date_to: datetime.date | None, db: Session, ) -> tuple[datetime.date, datetime.date, str]: """Résout (from, to, label) en dates locales America/Toronto.""" today = datetime.datetime.now(tz=TZ).date() if date_from and date_to: if date_from > date_to: raise HTTPException(status_code=422, detail="from doit précéder to") return date_from, date_to, f"du {date_from} au {date_to}" if period == "auj": return today, today, PERIOD_LABELS["auj"] if period == "annee": return datetime.date(today.year, 1, 1), today, PERIOD_LABELS["annee"] if period == "tout": firsts = [ db.execute(select(func.min(ApiRequest.ts))).scalar(), db.execute(select(func.min(CollectionRun.started_at))).scalar(), ] dates = [] for f in firsts: if f is not None: if f.tzinfo is None: f = f.replace(tzinfo=datetime.UTC) dates.append(f.astimezone(TZ).date()) return (min(dates) if dates else today), today, PERIOD_LABELS["tout"] if period in PERIOD_DAYS: return ( today - datetime.timedelta(days=PERIOD_DAYS[period] - 1), today, PERIOD_LABELS[period], ) raise HTTPException( status_code=422, detail=f"Période invalide : {period}. Valides : {', '.join(PERIOD_LABELS)}", ) def _utc_bounds( d_from: datetime.date, d_to: datetime.date ) -> tuple[datetime.datetime, datetime.datetime]: """Bornes UTC [00:00 from, 24:00 to] exprimées depuis les dates locales.""" start = datetime.datetime.combine(d_from, datetime.time.min, tzinfo=TZ) end = datetime.datetime.combine( d_to + datetime.timedelta(days=1), datetime.time.min, tzinfo=TZ ) return start.astimezone(datetime.UTC), end.astimezone(datetime.UTC) # ---------------------------------------------------------------- agrégats # Endpoints « bruit » exclus des métriques du tableau de bord : chemins # inconnus sondés par les scanners/bots (404, .env, etc.). Ils restent # journalisés dans api_requests pour la visibilité sécurité, mais ne doivent # pas polluer le taux d'erreur qui mesure la santé réelle de l'API. NOISE_ENDPOINTS = ("(autre)", "/api/v1/(autre)") def _fetch_requests( db: Session, d_from: datetime.date, d_to: datetime.date ) -> list[tuple[datetime.datetime, str, str, int, float]]: """Requêtes API de la fenêtre : (ts local, méthode, endpoint, statut, ms).""" lo, hi = _utc_bounds(d_from, d_to) rows = db.execute( select( ApiRequest.ts, ApiRequest.method, ApiRequest.endpoint, ApiRequest.status, ApiRequest.duration_ms, ).where( ApiRequest.ts >= lo, ApiRequest.ts < hi, ApiRequest.endpoint.notin_(NOISE_ENDPOINTS), ) ).all() out = [] for ts, method, endpoint, status, duration in rows: if ts.tzinfo is None: ts = ts.replace(tzinfo=datetime.UTC) out.append((ts.astimezone(TZ), method or "GET", endpoint, status, duration or 0.0)) return out def _fetch_runs( db: Session, d_from: datetime.date, d_to: datetime.date ) -> list[CollectionRun]: """Runs de collecte de la fenêtre (par date logique ``date_key``).""" return ( db.execute( select(CollectionRun) .where(CollectionRun.date_key >= d_from, CollectionRun.date_key <= d_to) .order_by(CollectionRun.started_at.desc()) ) .scalars() .all() ) def _fetch_health(db: Session) -> list[ConnectorHealth]: """État courant des connecteurs supervisés (hors ligne virtuelle _app).""" return ( db.execute( select(ConnectorHealth) .where(ConnectorHealth.source != "_app") .order_by(ConnectorHealth.service, ConnectorHealth.source) ) .scalars() .all() ) def _p95(values: list[float]) -> float: if not values: return 0.0 vs = sorted(values) idx = max(0, int(round(0.95 * (len(vs) - 1)))) return vs[idx] def _delta(cur: float, prev: float | None) -> float | None: if prev is None or prev == 0: return None return round((cur - prev) / prev * 100, 1) def _daterange(d_from: datetime.date, d_to: datetime.date) -> list[datetime.date]: n = (d_to - d_from).days + 1 return [d_from + datetime.timedelta(days=i) for i in range(n)] def _fr_int(n: float) -> str: return f"{int(n):,}".replace(",", " ") def _fr_num(n: float) -> str: if isinstance(n, float) and not float(n).is_integer(): return f"{n:,.1f}".replace(",", " ").replace(".", ",") return _fr_int(n) def _spark(points: list[dict], max_pts: int = 40) -> list[dict]: """Sous-échantillonne une série pour la sparkline d'un KPI (≤ max_pts).""" if len(points) <= max_pts: return points step = len(points) / max_pts out = [points[min(len(points) - 1, int(i * step))] for i in range(max_pts)] if out[-1] is not points[-1]: out[-1] = points[-1] return out def _hist(values: list[float], bins) -> list[dict]: out = [] for lo, hi, label in bins: c = sum(1 for v in values if lo <= v < hi) out.append({"label": label, "value": c}) while out and out[-1]["value"] == 0: out.pop() return out # ---------------------------------------------------------------- dashboard def _build_dashboard(db: Session, period: str, d_from, d_to, label) -> dict[str, Any]: """Calcule le JSON complet du contrat SPEC v2 depuis les données réelles.""" hourly = period == "auj" or d_from == d_to span = (d_to - d_from).days + 1 prev_to = d_from - datetime.timedelta(days=1) prev_from = prev_to - datetime.timedelta(days=span - 1) reqs = _fetch_requests(db, d_from, d_to) prev_reqs = _fetch_requests(db, prev_from, prev_to) runs = _fetch_runs(db, d_from, d_to) prev_runs = _fetch_runs(db, prev_from, prev_to) health = _fetch_health(db) days = _daterange(d_from, d_to) # -------- seaux temporels (jour, ou heure quand la fenêtre = 1 jour) if hourly: bucket_keys: list = list(range(24)) bucket_labels = [f"{h:02d}h" for h in range(24)] req_bucket = {i: [] for i in bucket_keys} for r in reqs: req_bucket[r[0].hour].append(r) else: bucket_keys = days bucket_labels = [d.isoformat() for d in days] req_bucket = {d: [] for d in days} for r in reqs: d = r[0].date() if d in req_bucket: req_bucket[d].append(r) def per_bucket(fn) -> list[dict]: return [ {"t": bucket_labels[i], "v": fn(req_bucket[k])} for i, k in enumerate(bucket_keys) ] calls_pts = per_bucket(len) err_pts = per_bucket(lambda rs: sum(1 for r in rs if r[3] >= 400)) avg_pts = per_bucket( lambda rs: round(sum(r[4] for r in rs) / len(rs), 1) if rs else 0 ) p95_pts = per_bucket(lambda rs: round(_p95([r[4] for r in rs]), 1) if rs else 0) errate_pts = per_bucket( lambda rs: round(sum(1 for r in rs if r[3] >= 400) / len(rs) * 100, 1) if rs else 0 ) # -------- agrégats requêtes API (période courante) total_calls = len(reqs) durations = [r[4] for r in reqs] avg_ms = round(sum(durations) / total_calls, 2) if total_calls else 0.0 p95_ms = round(_p95(durations), 2) errors = sum(1 for r in reqs if r[3] >= 400) err5xx = sum(1 for r in reqs if r[3] >= 500) err_rate = round(errors / total_calls * 100, 2) if total_calls else 0.0 success = total_calls - errors fast200 = sum(1 for d in durations if d < 200) active_endpoints = len({r[2] for r in reqs}) active_days = sum(1 for k in bucket_keys if req_bucket[k]) if not hourly else None # -------- agrégats requêtes API (période précédente, pour les deltas) prev_calls = len(prev_reqs) prev_durs = [r[4] for r in prev_reqs] prev_avg = (sum(prev_durs) / prev_calls) if prev_calls else None prev_p95 = _p95(prev_durs) if prev_calls else None prev_errors = sum(1 for r in prev_reqs if r[3] >= 400) prev_err5xx = sum(1 for r in prev_reqs if r[3] >= 500) prev_err = (prev_errors / prev_calls * 100) if prev_calls else None prev_endpoints = len({r[2] for r in prev_reqs}) # -------- agrégats runs de collecte ok_runs = [r for r in runs if r.status in ("success", "retried")] failed_runs = [r for r in runs if r.status == "failed"] records_total = sum(r.records_count for r in ok_runs) active_services = len({r.service for r in runs}) prev_ok = [r for r in prev_runs if r.status in ("success", "retried")] prev_failed = [r for r in prev_runs if r.status == "failed"] prev_records = sum(r.records_count for r in prev_ok) rec_day: dict[datetime.date, int] = {d: 0 for d in days} okrun_day: dict[datetime.date, int] = {d: 0 for d in days} for r in ok_runs: if r.date_key in rec_day: rec_day[r.date_key] += r.records_count okrun_day[r.date_key] += 1 # -------- KPI (uniquement des mesures réelles) + sparklines kpis: list[dict[str, Any]] = [] def kpi(id_, label_, value, unit="", delta=None, invert=False, spark=None): d: dict[str, Any] = {"id": id_, "label": label_, "value": value, "unit": unit} if delta is not None: d["delta_pct"] = delta good = (delta <= 0) if invert else (delta >= 0) d["direction"] = "up" if good else "down" if invert: d["invert"] = True if spark and len(spark) > 1: d["spark"] = _spark(spark) kpis.append(d) kpi( "calls", "Appels API", total_calls, "", _delta(total_calls, prev_calls), spark=calls_pts, ) if not hourly and span > 1: kpi( "calls_day", "Appels par jour (moyenne)", round(total_calls / span, 1), "", _delta(total_calls / span, prev_calls / span if prev_calls else None), ) kpi( "latency", "Latence moyenne", avg_ms, "ms", _delta(avg_ms, round(prev_avg, 2) if prev_avg else None), invert=True, spark=avg_pts, ) kpi( "p95", "Latence p95", p95_ms, "ms", _delta(p95_ms, round(prev_p95, 2) if prev_p95 else None), invert=True, spark=p95_pts, ) kpi( "errors", "Taux d'erreur (HTTP ≥ 400)", err_rate, "%", _delta(err_rate, round(prev_err, 2) if prev_err else None), invert=True, spark=errate_pts, ) kpi( "err5xx", "Erreurs serveur (5xx)", err5xx, "", _delta(err5xx, prev_err5xx or None), invert=True, ) kpi( "endpoints", "Endpoints actifs", active_endpoints, "", _delta(active_endpoints, prev_endpoints or None), ) kpi( "runs_ok", "Collectes réussies", len(ok_runs), "", _delta(len(ok_runs), len(prev_ok) or None), spark=None if hourly else [{"t": d.isoformat(), "v": okrun_day[d]} for d in days], ) kpi( "runs_failed", "Collectes échouées", len(failed_runs), "", _delta(len(failed_runs), len(prev_failed) or None), invert=True, ) kpi( "records", "Enregistrements collectés", records_total, "", _delta(records_total, prev_records or None), spark=None if hourly else [{"t": d.isoformat(), "v": rec_day[d]} for d in days], ) if active_services: kpi("services", "Services KA collectés", active_services) # -------- jauges : taux & couvertures (données réelles seulement) gauges: list[dict[str, Any]] = [] if total_calls: gauges.append( { "id": "success", "label": "Taux de succès HTTP (statut < 400)", "value": round(success / total_calls * 100, 1), "max": 100, "unit": "%", } ) gauges.append( { "id": "fast", "label": "Réponses servies en moins de 200 ms", "value": round(fast200 / total_calls * 100, 1), "max": 100, "unit": "%", } ) if runs: gauges.append( { "id": "runs", "label": "Taux de réussite des collectes", "value": round(len(ok_runs) / len(runs) * 100, 1), "max": 100, "unit": "%", } ) if health: ok_h = sum(1 for h in health if h.status == "ok") gauges.append( { "id": "connectors", "label": f"Connecteurs de l'écosystème en santé ({ok_h}/{len(health)})", "value": round(ok_h / len(health) * 100, 1), "max": 100, "unit": "%", } ) if active_days is not None and span > 1 and total_calls: gauges.append( { "id": "activedays", "label": f"Jours avec trafic API journalisé ({active_days}/{span})", "value": round(active_days / span * 100, 1), "max": 100, "unit": "%", } ) # -------- séries temporelles series: list[dict[str, Any]] = [] suffix = "par heure" if hourly else "par jour" calls_serie: dict[str, Any] = { "id": "calls", "title": f"Appels API {suffix}", "unit": "appels", "kind": "line", "points": calls_pts, } if prev_reqs and not hourly: prev_days = _daterange(prev_from, prev_to) prev_by = {d: 0 for d in prev_days} for r in prev_reqs: d = r[0].date() if d in prev_by: prev_by[d] += 1 calls_serie["compare"] = [ {"t": d.isoformat(), "v": prev_by[d]} for d in prev_days ] if reqs: series.append(calls_serie) series.append( { "id": "errors", "title": f"Erreurs HTTP (≥ 400) {suffix}", "unit": "erreurs", "kind": "bar", "points": err_pts, } ) series.append( { "id": "latency", "title": f"Latence moyenne {suffix} (ms)", "unit": "ms", "kind": "line", "points": avg_pts, } ) if runs and not hourly: series.append( { "id": "records", "title": "Enregistrements collectés par jour", "unit": "enregistrements", "kind": "line", "points": [{"t": d.isoformat(), "v": rec_day[d]} for d in days], } ) # -------- multi-séries (≤ 4 séries chacune) multiseries: list[dict[str, Any]] = [] if reqs and len(bucket_keys) > 1: multiseries.append( { "id": "lat", "title": f"Latence {suffix} — moyenne vs p95 (ms)", "unit": "ms", "series": [ {"label": "Moyenne", "points": avg_pts}, {"label": "p95", "points": p95_pts}, ], } ) ep_stats: dict[str, dict[str, Any]] = {} for _, _, endpoint, status, dur in reqs: s = ep_stats.setdefault(endpoint, {"calls": 0, "durs": [], "errors": 0}) s["calls"] += 1 s["durs"].append(dur) if status >= 400: s["errors"] += 1 top_eps = sorted(ep_stats.items(), key=lambda kv: kv[1]["calls"], reverse=True) if len(top_eps) >= 2 and len(bucket_keys) > 1: top4 = [ep for ep, _ in top_eps[:4]] ep_bucket = { ep: {k: 0 for k in bucket_keys} for ep in top4 } for r in reqs: if r[2] in ep_bucket: key = r[0].hour if hourly else r[0].date() if key in ep_bucket[r[2]]: ep_bucket[r[2]][key] += 1 multiseries.append( { "id": "top_eps", "title": f"Appels {suffix} — top {len(top4)} endpoints", "unit": "appels", "series": [ { "label": ep, "points": [ {"t": bucket_labels[i], "v": ep_bucket[ep][k]} for i, k in enumerate(bucket_keys) ], } for ep in top4 ], } ) # -------- barres empilées : composition dans le temps stacked: list[dict[str, Any]] = [] classes = [("2xx", 200, 300), ("3xx", 300, 400), ("4xx", 400, 500), ("5xx", 500, 600)] if reqs: present = [ (name, lo, hi) for name, lo, hi in classes if any(lo <= r[3] < hi for r in reqs) ] if present: stacked.append( { "id": "status", "title": f"Appels {suffix} par classe de statut HTTP", "unit": "appels", "keys": [c[0] for c in present], "points": [ { "t": bucket_labels[i], "values": [ sum(1 for r in req_bucket[k] if lo <= r[3] < hi) for _, lo, hi in present ], } for i, k in enumerate(bucket_keys) ], } ) svc_records: dict[str, int] = {} for r in ok_runs: svc_records[r.service] = svc_records.get(r.service, 0) + r.records_count if svc_records and not hourly and len(days) > 1: top_svcs = [ s for s, _ in sorted(svc_records.items(), key=lambda kv: kv[1], reverse=True) ][:6] svc_day = {s: {d: 0 for d in days} for s in top_svcs} for r in ok_runs: if r.service in svc_day and r.date_key in svc_day[r.service]: svc_day[r.service][r.date_key] += r.records_count stacked.append( { "id": "services_day", "title": "Enregistrements collectés par jour et par service", "unit": "enregistrements", "keys": [SERVICE_NAMES.get(s, s) for s in top_svcs], "points": [ { "t": d.isoformat(), "values": [svc_day[s][d] for s in top_svcs], } for d in days ], } ) # -------- répartitions (avec deltas honnêtes vs période précédente) prev_ep_calls: dict[str, int] = {} prev_status_cls: dict[str, int] = {} prev_methods: dict[str, int] = {} for _, method, endpoint, status, _dur in prev_reqs: prev_ep_calls[endpoint] = prev_ep_calls.get(endpoint, 0) + 1 cls = f"{status // 100}xx" prev_status_cls[cls] = prev_status_cls.get(cls, 0) + 1 prev_methods[method] = prev_methods.get(method, 0) + 1 breakdowns: list[dict[str, Any]] = [] if top_eps: items = [] for ep, s in top_eps[:12]: it = {"label": ep, "value": s["calls"]} d = _delta(s["calls"], prev_ep_calls.get(ep)) if d is not None: it["delta_pct"] = d items.append(it) breakdowns.append( { "id": "top_endpoints", "title": "Top endpoints (appels)", "kind": "bars", "items": items, } ) if reqs: status_counts: dict[int, int] = {} for r in reqs: status_counts[r[3]] = status_counts.get(r[3], 0) + 1 cls_counts: dict[str, int] = {} for st, c in status_counts.items(): cls = f"{st // 100}xx" cls_counts[cls] = cls_counts.get(cls, 0) + c items = [] for cls in sorted(cls_counts): it = {"label": cls, "value": cls_counts[cls]} d = _delta(cls_counts[cls], prev_status_cls.get(cls)) if d is not None: it["delta_pct"] = d items.append(it) breakdowns.append( { "id": "status", "title": "Appels par classe de statut HTTP", "kind": "donut", "items": items, } ) meth_counts: dict[str, int] = {} for r in reqs: meth_counts[r[1]] = meth_counts.get(r[1], 0) + 1 items = [] for m, c in sorted(meth_counts.items(), key=lambda kv: kv[1], reverse=True): it = {"label": m, "value": c} d = _delta(c, prev_methods.get(m)) if d is not None: it["delta_pct"] = d items.append(it) breakdowns.append( { "id": "methods", "title": "Appels par méthode HTTP", "kind": "bars", "items": items, } ) if svc_records: prev_svc: dict[str, int] = {} for r in prev_ok: prev_svc[r.service] = prev_svc.get(r.service, 0) + r.records_count items = [] for s, v in sorted(svc_records.items(), key=lambda kv: kv[1], reverse=True): it = {"label": SERVICE_NAMES.get(s, s), "value": v} d = _delta(v, prev_svc.get(s)) if d is not None: it["delta_pct"] = d items.append(it) breakdowns.append( { "id": "services", "title": "Enregistrements collectés par service", "kind": "donut", "items": items, } ) # -------- distributions (histogrammes) distributions: list[dict[str, Any]] = [] if durations: bins = _hist(durations, LAT_BINS) if bins: distributions.append( { "id": "latency_hist", "title": "Distribution des latences (temps de réponse)", "unit": "appels", "bins": bins, } ) run_durs = [r.duration_seconds for r in runs] if run_durs: bins = _hist(run_durs, DUR_BINS) if bins: distributions.append( { "id": "rundur_hist", "title": "Distribution des durées de collecte", "unit": "runs", "bins": bins, } ) # -------- heatmap calendrier : appels API par jour heatmap = None if reqs and not hourly: cells = [ {"date": k.isoformat(), "value": len(req_bucket[k])} for k in bucket_keys if req_bucket[k] ] if cells: heatmap = {"title": "Appels API par jour", "cells": cells} # -------- heatmap horaire 7×24 : appels par jour de semaine × heure hourly_map = None if reqs: grid: dict[tuple[int, int], int] = {} for r in reqs: key = (r[0].weekday(), r[0].hour) # 0 = lundi (contrat SPEC) grid[key] = grid.get(key, 0) + 1 if grid: hourly_map = { "title": "Appels API par jour de semaine et heure", "cells": [ {"dow": dw, "hour": h, "value": v} for (dw, h), v in sorted(grid.items()) ], } # -------- tableaux tables: list[dict[str, Any]] = [] if top_eps: tables.append( { "id": "endpoints", "title": "Endpoints — appels, latence et erreurs", "columns": [ "Endpoint", "Appels", "Latence moy. (ms)", "p95 (ms)", "Erreurs", "Taux d'erreur", ], "rows": [ [ ep, s["calls"], round(sum(s["durs"]) / len(s["durs"]), 1), round(_p95(s["durs"]), 1), s["errors"], f"{s['errors'] / s['calls'] * 100:.1f} %".replace(".", ","), ] for ep, s in top_eps ], } ) if reqs and len(bucket_keys) > 1: rows = [] for i, k in enumerate(bucket_keys): rs = req_bucket[k] if not rs and hourly: continue durs = [r[4] for r in rs] rows.append( [ bucket_labels[i], len(rs), sum(1 for r in rs if r[3] >= 400), round(sum(durs) / len(durs), 1) if durs else 0, round(_p95(durs), 1) if durs else 0, ] ) tables.append( { "id": "daily", "title": "Activité par heure" if hourly else "Activité par jour", "columns": [ "Heure" if hourly else "Date", "Appels", "Erreurs", "Latence moy. (ms)", "p95 (ms)", ], "rows": rows, } ) if runs: svc_agg: dict[str, dict[str, Any]] = {} for r in runs: a = svc_agg.setdefault( r.service, {"runs": 0, "ok": 0, "failed": 0, "records": 0, "durs": []} ) a["runs"] += 1 if r.status in ("success", "retried"): a["ok"] += 1 a["records"] += r.records_count else: a["failed"] += 1 a["durs"].append(r.duration_seconds) tables.append( { "id": "services", "title": "Services KA — collectes de la période", "columns": [ "Service", "Runs", "Réussis", "Échoués", "Enregistrements", "Durée moy. (s)", ], "rows": [ [ SERVICE_NAMES.get(s, s), a["runs"], a["ok"], a["failed"], a["records"], round(sum(a["durs"]) / len(a["durs"]), 1), ] for s, a in sorted( svc_agg.items(), key=lambda kv: kv[1]["records"], reverse=True ) ], } ) tables.append( { "id": "runs", "title": "Derniers runs de collecte", "columns": ["Service", "Date", "Statut", "Enregistrements", "Durée (s)"], "rows": [ [ SERVICE_NAMES.get(r.service, r.service), r.date_key.isoformat(), RUN_STATUS_FR.get(r.status, r.status), r.records_count, round(r.duration_seconds, 1), ] for r in runs[:80] ], } ) if health: tables.append( { "id": "connectors", "title": "État des connecteurs supervisés (temps réel)", "columns": [ "Service", "Source", "Statut", "Dernier succès", "Échecs consécutifs", ], "rows": [ [ SERVICE_NAMES.get(h.service, h.service), h.source, HEALTH_FR.get(h.status, h.status), ( h.last_success.astimezone(TZ).strftime("%Y-%m-%d %H:%M") if h.last_success else "—" ), h.consecutive_failures, ] for h in sorted( health, key=lambda h: (h.status == "ok", h.service, h.source), ) ], } ) # -------- records & faits marquants (6-12, générés depuis les données) records: list[dict[str, Any]] = [] if reqs and not hourly: best_day = max(bucket_keys, key=lambda k: len(req_bucket[k])) if req_bucket[best_day]: records.append( { "label": "Jour record d'appels API", "value": f"{_fr_int(len(req_bucket[best_day]))} appels", "date": best_day.isoformat(), } ) if hourly_map: top_cell = max(hourly_map["cells"], key=lambda c: c["value"]) dows = ["lundi", "mardi", "mercredi", "jeudi", "vendredi", "samedi", "dimanche"] records.append( { "label": "Créneau horaire le plus actif", "value": ( f"{dows[top_cell['dow']]} {top_cell['hour']:02d}h — " f"{_fr_int(top_cell['value'])} appels" ), } ) if top_eps: records.append( { "label": "Endpoint le plus sollicité", "value": f"{top_eps[0][0]} — {_fr_int(top_eps[0][1]['calls'])} appels", } ) grow = [ (ep, _delta(s["calls"], prev_ep_calls.get(ep))) for ep, s in top_eps if prev_ep_calls.get(ep, 0) >= 5 ] grow = [(ep, d) for ep, d in grow if d is not None and d > 0] if grow: ep, d = max(grow, key=lambda kv: kv[1]) records.append( { "label": "Plus forte croissance d'endpoint", "value": f"{ep} — +{_fr_num(d)} %", } ) if durations: records.append( { "label": "Réponse la plus rapide de la période", "value": f"{_fr_num(round(min(durations), 1))} ms", } ) records.append( { "label": "Réponse la plus lente de la période", "value": f"{_fr_num(round(max(durations), 1))} ms", } ) if reqs and not hourly: clean_days = [ k for k in bucket_keys if req_bucket[k] and not any(r[3] >= 400 for r in req_bucket[k]) ] if clean_days: best = max(clean_days, key=lambda k: len(req_bucket[k])) records.append( { "label": "Meilleure journée sans erreur HTTP", "value": f"{_fr_int(len(req_bucket[best]))} appels, 0 erreur", "date": best.isoformat(), } ) if ok_runs: biggest = max(ok_runs, key=lambda r: r.records_count) records.append( { "label": "Run de collecte le plus volumineux", "value": ( f"{_fr_int(biggest.records_count)} enregistrements " f"({SERVICE_NAMES.get(biggest.service, biggest.service)})" ), "date": biggest.date_key.isoformat(), } ) fastest = min(ok_runs, key=lambda r: r.duration_seconds) records.append( { "label": "Collecte la plus rapide", "value": ( f"{fastest.duration_seconds:.1f} s " f"({SERVICE_NAMES.get(fastest.service, fastest.service)})" ).replace(".", ","), "date": fastest.date_key.isoformat(), } ) total_logged = db.execute(select(func.count(ApiRequest.id))).scalar() or 0 if total_logged: records.append( { "label": "Appels journalisés au total (90 jours de rétention)", "value": f"{_fr_int(total_logged)} appels", } ) # -------- couverture de mesure (depuis quand les appels sont journalisés) first_req = db.execute(select(func.min(ApiRequest.ts))).scalar() if first_req is not None and first_req.tzinfo is None: first_req = first_req.replace(tzinfo=datetime.UTC) dash: dict[str, Any] = { "updated": datetime.datetime.now(tz=TZ).isoformat(timespec="seconds"), "period": {"from": d_from.isoformat(), "to": d_to.isoformat(), "label": label}, "kpis": kpis, "series": series, "breakdowns": breakdowns, "tables": tables, "records": records, "coverage": { "api_requests_since": ( first_req.astimezone(TZ).isoformat(timespec="seconds") if first_req else None ) }, } if gauges: dash["gauges"] = gauges if multiseries: dash["multiseries"] = multiseries if stacked: dash["stacked"] = stacked if distributions: dash["distributions"] = distributions if heatmap: dash["heatmap"] = heatmap if hourly_map: dash["hourly"] = hourly_map return dash def _dashboard_cached( db: Session, period: str, date_from: datetime.date | None, date_to: datetime.date | None, ) -> dict[str, Any]: d_from, d_to, label = _resolve_period(period, date_from, date_to, db) key = (period, d_from.isoformat(), d_to.isoformat()) now = time.monotonic() with _cache_lock: hit = _cache.get(key) if hit and now - hit[0] < CACHE_TTL: return hit[1] dash = _build_dashboard(db, period, d_from, d_to, label) with _cache_lock: if len(_cache) > 64: _cache.clear() _cache[key] = (now, dash) return dash def _normalize_mode(mode: str | None) -> str: """SPEC v2 : un mode inconnu ou absent retombe sur ``complet``.""" m = (mode or "").strip().lower() return m if m in kapdf.REPORT_MODES else "complet" # ---------------------------------------------------------------- routes @router.get("") def stats_summary(db: Session = Depends(get_db)) -> dict[str, Any]: """Résumé compact de la plateforme (contrat léger /api/stats commun KA). Version allégée du tableau de bord pour les moniteurs et les apps sœurs : activité API 7 jours, runs de collecte, dernier run par service et santé des connecteurs supervisés. """ today = datetime.datetime.now(TZ).date() d_from = today - datetime.timedelta(days=6) reqs = _fetch_requests(db, d_from, today) runs = _fetch_runs(db, d_from, today) health = _fetch_health(db) errors = sum(1 for r in reqs if r[3] >= 400) ok_runs = [r for r in runs if r.status in ("success", "retried")] last_run: dict[str, Any] = {} for r in runs: # triés par started_at décroissant if r.service not in last_run: last_run[r.service] = { "date": r.date_key.isoformat(), "status": r.status, "records": r.records_count, } connectors = dict.fromkeys(("ok", "degraded", "broken", "stale"), 0) for h in health: if h.status in connectors: connectors[h.status] += 1 return envelope( { "platform": "API-KA", "period": {"from": d_from.isoformat(), "to": today.isoformat()}, "api": { "calls_7d": len(reqs), "error_rate_7d": round(errors / len(reqs) * 100, 2) if reqs else 0.0, }, "runs_7d": { "total": len(runs), "ok": len(ok_runs), "failed": sum(1 for r in runs if r.status == "failed"), "records": sum(r.records_count for r in ok_runs), }, "last_run": last_run, "connectors": connectors, } ) @router.get("/dashboard") def stats_dashboard( period: str = Query("30j", description="auj, 7j, 30j, 3m, 6m, 12m, annee, tout"), date_from: datetime.date | None = Query(None, alias="from"), date_to: datetime.date | None = Query(None, alias="to"), db: Session = Depends(get_db), ) -> dict[str, Any]: """Tableau de bord analytique de la plateforme (contrat commun Groupe KA).""" return envelope(_dashboard_cached(db, period, date_from, date_to)) @router.get("/report") def stats_report( period: str = Query("30j", description="auj, 7j, 30j, 3m, 6m, 12m, annee, tout"), date_from: datetime.date | None = Query(None, alias="from"), date_to: datetime.date | None = Query(None, alias="to"), mode: str = Query( "complet", description="complet, synthese, tendances, repartitions ou donnees", ), db: Session = Depends(get_db), ) -> Response: """Rapport statistique PDF estampillé Groupe-KA (5 rapports, SPEC v2).""" mode = _normalize_mode(mode) dash = _dashboard_cached(db, period, date_from, date_to) pdf_bytes = kapdf.GroupeKAReport(site=SITE, dashboard=dash, mode=mode).build() fname = kapdf.filename(PLATFORM_ID, period, mode) return Response( content=pdf_bytes, media_type="application/pdf", headers={"Content-Disposition": f'attachment; filename="{fname}"'}, ) @router.get("/catalog") def stats_catalog( period: str = Query("30j", description="auj, 7j, 30j, 3m, 6m, 12m, annee, tout"), date_from: datetime.date | None = Query(None, alias="from"), date_to: datetime.date | None = Query(None, alias="to"), db: Session = Depends(get_db), ) -> dict[str, Any]: """v3 — blocs composables pour le constructeur de rapports personnalisés.""" dash = _dashboard_cached(db, period, date_from, date_to) return { "updated": dash.get("updated"), "period": dash.get("period"), "blocks": kapdf.catalog(dash), } @router.post("/report/custom") def stats_report_custom( spec: dict = Body(...), db: Session = Depends(get_db), ) -> Response: """v3 — rapport PDF personnalisé : ``{"title", "period", "from", "to", "blocks": [{"key": "series:…", "render": "bar"}, …]}`` (SPEC.md §3bis).""" period = str(spec.get("period") or "30j") if period not in PERIOD_LABELS: period = "30j" def _date(v: Any) -> datetime.date | None: try: return datetime.date.fromisoformat(str(v)) if v else None except ValueError: return None dash = _dashboard_cached(db, period, _date(spec.get("from")), _date(spec.get("to"))) known = {b["key"] for b in kapdf.catalog(dash)} blocks = [ b for b in (spec.get("blocks") or []) if isinstance(b, dict) and b.get("key") in known ][:40] if not blocks: raise HTTPException(400, "Aucun bloc valide dans la composition") pdf_bytes = kapdf.GroupeKAReport( site=SITE, dashboard=dash, mode=kapdf.CUSTOM_MODE, spec={"title": str(spec.get("title") or "")[:80], "blocks": blocks}, ).build() fname = kapdf.filename(PLATFORM_ID, period, kapdf.CUSTOM_MODE) return Response( content=pdf_bytes, media_type="application/pdf", headers={"Content-Disposition": f'attachment; filename="{fname}"'}, ) # ------------------------------------------------- rapport écosystème (PDF) async def _fetch_satellite_dashboard( client: httpx.AsyncClient, site: dict, period: str ) -> dict[str, Any]: """Dashboard d'une plateforme sœur — jamais d'exception : une plateforme injoignable est retournée avec ``dashboard=None`` et un motif lisible.""" url = f"https://{site['domain']}/api/stats/dashboard" try: resp = await client.get(url, params={"period": period}) resp.raise_for_status() payload = resp.json() # certaines plateformes enveloppent la réponse ({success, data, meta}) if ( isinstance(payload, dict) and "kpis" not in payload and isinstance(payload.get("data"), dict) ): payload = payload["data"] if not isinstance(payload, dict) or not payload.get("kpis"): raise ValueError("réponse hors contrat (kpis manquants)") return {"site": site, "dashboard": payload, "error": None} except httpx.TimeoutException: return {"site": site, "dashboard": None, "error": "délai dépassé (15 s)"} except httpx.HTTPStatusError as exc: return { "site": site, "dashboard": None, "error": f"HTTP {exc.response.status_code}", } except Exception as exc: # noqa: BLE001 — motif affiché dans le PDF return { "site": site, "dashboard": None, "error": (str(exc) or exc.__class__.__name__)[:80], } @router.get("/ecosystem-report") async def stats_ecosystem_report( period: str = Query("30j", description="auj, 7j, 30j, 3m, 6m, 12m, annee, tout"), mode: str | None = Query( None, description=( "complet, synthese, tendances, repartitions ou donnees " "(optionnel — inconnu ou absent = complet)" ), ), db: Session = Depends(get_db), ) -> Response: """Rapport écosystème Groupe-KA : UN SEUL PDF consolidant les 12 plateformes. Les dashboards des plateformes sœurs sont récupérés en parallèle (httpx, délai 15 s), celui d'API-KA est calculé en local. Le ``mode`` (5 rapports du contrat v2) est passé au moteur ; un mode inconnu ou absent retombe sur ``complet``. Cache 10 min par (période, mode).""" mode = _normalize_mode(mode) if period not in PERIOD_LABELS: raise HTTPException( status_code=422, detail=f"Période invalide : {period}. Valides : {', '.join(PERIOD_LABELS)}", ) fname = kapdf.filename("ecosysteme", period, mode) def _pdf_response(pdf_bytes: bytes, fname: str) -> Response: return Response( content=pdf_bytes, media_type="application/pdf", headers={"Content-Disposition": f'attachment; filename="{fname}"'}, ) key = (period, mode) now = time.monotonic() with _eco_lock: hit = _eco_cache.get(key) if hit and now - hit[0] < ECO_CACHE_TTL: return _pdf_response(hit[1], hit[2]) satellites = [s for s in ECO_SITES if s["id"] != PLATFORM_ID] own_site = next(s for s in ECO_SITES if s["id"] == PLATFORM_ID) async with httpx.AsyncClient( timeout=ECO_FETCH_TIMEOUT, follow_redirects=True, headers={ "accept": "application/json", "user-agent": "apika-ecosystem-report/1.0 (+https://www.api-ka.com)", }, ) as client: results = await asyncio.gather( *(_fetch_satellite_dashboard(client, s, period) for s in satellites) ) def _build() -> bytes: try: own = { "site": own_site, "dashboard": _dashboard_cached(db, period, None, None), "error": None, } except Exception as exc: # noqa: BLE001 own = {"site": own_site, "dashboard": None, "error": str(exc)[:80]} platforms = [*results, own] eco_period = next( ( pl["dashboard"]["period"] for pl in platforms if pl["dashboard"] and pl["dashboard"].get("period") ), {"label": PERIOD_LABELS[period]}, ) return ecopdf.EcosystemReport( platforms=platforms, period=eco_period, mode=mode ).build() pdf_bytes = await asyncio.to_thread(_build) with _eco_lock: if len(_eco_cache) > 32: _eco_cache.clear() _eco_cache[key] = (now, pdf_bytes, fname) return _pdf_response(pdf_bytes, fname)