# Trouve-KA — module Stats : dashboard analytique + rapport PDF Groupe-KA # Author: Simon-Pierre Boucher # Contact: contact@spboucher.ai """Statistiques du moteur (contrat ka-ui/stats/SPEC.md v2). - `dashboard()` agrège Postgres (documents, crawl_attempts, frontier_items, domains, search_queries) en JSON conforme au contrat /api/stats/dashboard : KPI (+sparklines), jauges, séries, multi-séries, empilées, répartitions, distributions, heatmap calendrier, heatmap horaire 7×24, tableaux, records. - Cache mémoire par période (>= 5 min) : les tables sont volumineuses (crawl_attempts ~5 M lignes), on ne re-agrège jamais à chaque appel. Toutes les requêtes sont bornées sur des colonnes indexées (fetched_at, first_indexed_at, created_at) — jamais de scan complet par requête HTTP. Les agrégats globaux lourds (profil documents, profondeurs frontier) vivent dans un cache dédié (10 min à 1 h). - AUCUNE donnée inventée : une section sans données est simplement absente (le front affiche « Pas encore mesuré »). - Purge opportuniste des journaux de recherche (> 90 jours), 1×/jour max. """ from __future__ import annotations import asyncio import time from datetime import datetime, timedelta from zoneinfo import ZoneInfo from trouveka.logging import get_logger log = get_logger("stats") TZ = ZoneInfo("America/Toronto") SITE = { "wordmark": "Trouve·Ka", "domain": "www.trouve-ka.com", "accent": "#1c7ed6", "tagline": "Le moteur de recherche du Groupe 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", } OUTCOME_LABELS = { "indexed": "Indexée", "duplicate": "Doublon", "error": "Erreur", "not_quebec": "Hors périmètre", "robots_blocked": "Bloquée (robots.txt)", "redirect": "Redirection", "unchanged": "Inchangée", } FRONTIER_LABELS = { "pending": "En attente", "in_progress": "En cours", "done": "Complétée", "failed": "Échouée", "blocked": "Bloquée", } LANG_LABELS = {"fr": "Français", "en": "Anglais", "?": "Non détectée"} HTTP_CLASSES = ["2xx", "3xx", "4xx", "5xx", "Sans réponse"] # Cache mémoire : clé période → (expiration monotonic, payload) _cache: dict[str, tuple[float, dict]] = {} _global_cache: dict[str, tuple[float, object]] = {} _build_lock = asyncio.Lock() _last_purge = 0.0 CACHE_TTL = 300 # 5 min par période CACHE_TTL_TOUT = 900 # « tout » agrège 5 M lignes : 15 min GLOBAL_TTL = 600 # heatmap / profil documents / records globaux DEEP_TTL = 3600 # profondeurs frontier (scan 3,7 M lignes) : 1 h def _fr_int(n: int | float) -> str: return f"{int(n):,}".replace(",", " ") def _delta(cur: float | None, prev: float | None) -> float | None: """Variation % vs période précédente. Pas de période de référence → None.""" if cur is None or prev is None or prev == 0: return None # pas de période de référence mesurée → pas de faux delta return round(100.0 * (cur - prev) / prev, 1) def _spark(points: list[dict], max_pts: int = 40) -> list[dict] | None: """Sous-échantillonne une série pour la sparkline d'un KPI (≥ 2 points).""" pts = [p for p in points if isinstance(p.get("v"), (int, float))] if len(pts) < 2: return None if len(pts) <= max_pts: return pts step = len(pts) / max_pts out = [pts[int(i * step)] for i in range(max_pts)] if out[-1] is not pts[-1]: out.append(pts[-1]) return out def _mb(nbytes: int) -> tuple[float, str]: """Octets → (valeur, unité) lisibles (Mo / Go / To).""" if nbytes >= 1e12: return round(nbytes / 1e12, 2), "To" if nbytes >= 1e9: return round(nbytes / 1e9, 2), "Go" return round(nbytes / 1e6, 1), "Mo" def resolve_period( period: str, from_s: str | None, to_s: str | None, dataset_start: datetime ) -> tuple[datetime, datetime, str, datetime | None, datetime | None]: """→ (start, end, label, prev_start, prev_end) en heure de l'Est.""" now = datetime.now(TZ) if from_s and to_s: start = datetime.strptime(from_s, "%Y-%m-%d").replace(tzinfo=TZ) end = datetime.strptime(to_s, "%Y-%m-%d").replace(tzinfo=TZ) + timedelta(days=1) end = min(end, now) label = f"du {from_s} au {to_s}" elif period == "auj": start = now.replace(hour=0, minute=0, second=0, microsecond=0) end, label = now, PERIOD_LABELS["auj"] elif period == "annee": start = now.replace(month=1, day=1, hour=0, minute=0, second=0, microsecond=0) end, label = now, PERIOD_LABELS["annee"] elif period == "tout": start = dataset_start.astimezone(TZ) end, label = now, PERIOD_LABELS["tout"] else: days = {"7j": 7, "30j": 30, "3m": 91, "6m": 182, "12m": 365}.get(period, 30) start = now - timedelta(days=days) end, label = now, PERIOD_LABELS.get(period, PERIOD_LABELS["30j"]) if start >= end: start = end - timedelta(days=1) if period == "tout" and not (from_s and to_s): return start, end, label, None, None span = end - start return start, end, label, start - span, start def _bucket_grid(start: datetime, end: datetime, bucket: str) -> tuple[list[datetime], list[str]]: """Clés (datetime naïfs, heure locale — même forme que date_trunc AT TIME ZONE) et libellés d'axe pour une plage, trous compris.""" keys: list[datetime] = [] labels: list[str] = [] if bucket == "hour": cur = start.replace(minute=0, second=0, microsecond=0) while cur < end: k = cur.astimezone(TZ).replace(tzinfo=None, minute=0, second=0, microsecond=0) keys.append(k) labels.append(k.strftime("%Hh")) cur += timedelta(hours=1) else: cur_d = start.astimezone(TZ).date() end_d = end.astimezone(TZ).date() while cur_d <= end_d: keys.append(datetime(cur_d.year, cur_d.month, cur_d.day)) labels.append(cur_d.isoformat()) cur_d += timedelta(days=1) return keys, labels async def _dataset_start(db) -> datetime: """Première trace du crawler (min borné par index) — cache 1 h.""" hit = _global_cache.get("dataset_start") if hit and hit[0] > time.monotonic(): return hit[1] # type: ignore[return-value] row = await db.pool.fetchrow("SELECT min(fetched_at) AS t FROM crawl_attempts") start = row["t"] or datetime.now(TZ) - timedelta(days=1) _global_cache["dataset_start"] = (time.monotonic() + 3600, start) return start async def _bucket_counts( db, table: str, ts_col: str, start: datetime, end: datetime, bucket: str, extra_where: str = "" ) -> list[dict]: """Comptes par jour (ou heure) en heure locale, trous remplis de zéros. Toujours borné sur la colonne temporelle indexée — jamais de scan complet. """ unit = "hour" if bucket == "hour" else "day" rows = await db.pool.fetch( f""" SELECT date_trunc('{unit}', {ts_col} AT TIME ZONE 'America/Toronto') AS b, count(*)::int AS n FROM {table} WHERE {ts_col} >= $1 AND {ts_col} < $2 {extra_where} GROUP BY 1 ORDER BY 1 """, start, end, ) got = {r["b"]: r["n"] for r in rows} keys, labels = _bucket_grid(start, end, bucket) return [{"t": lab, "v": got.get(k, 0)} for k, lab in zip(keys, labels)] async def _heatmap_cells(db) -> list[dict]: """Indexations/jour sur 182 jours — partagé entre périodes (cache 10 min).""" hit = _global_cache.get("heatmap") if hit and hit[0] > time.monotonic(): return hit[1] # type: ignore[return-value] since = datetime.now(TZ) - timedelta(days=182) rows = await db.pool.fetch( """ SELECT (first_indexed_at AT TIME ZONE 'America/Toronto')::date AS d, count(*)::int AS n FROM documents WHERE first_indexed_at >= $1 GROUP BY 1 ORDER BY 1 """, since, ) cells = [{"date": r["d"].isoformat(), "value": r["n"]} for r in rows] _global_cache["heatmap"] = (time.monotonic() + GLOBAL_TTL, cells) return cells async def _doc_profile(db) -> dict: """Profil global de l'index (1 seul passage sur documents, cache 10 min) : total, fraîcheur < 30 j, langues, score Québec moyen.""" hit = _global_cache.get("doc_profile") if hit and hit[0] > time.monotonic(): return hit[1] # type: ignore[return-value] fresh_since = datetime.now(TZ) - timedelta(days=30) rows = await db.pool.fetch( """ SELECT coalesce(nullif(language, ''), '?') AS lang, count(*)::int AS n, count(*) FILTER (WHERE last_indexed_at >= $1)::int AS fresh, sum(page_quebec_score)::float AS qs_sum FROM documents GROUP BY 1 """, fresh_since, ) total = sum(r["n"] for r in rows) prof = { "total": total, "fresh30": sum(r["fresh"] for r in rows), "avg_qs": round(sum(r["qs_sum"] or 0 for r in rows) / total, 2) if total else None, "langs": sorted( ({"label": LANG_LABELS.get(r["lang"], r["lang"]), "value": r["n"]} for r in rows), key=lambda x: -x["value"], ), } _global_cache["doc_profile"] = (time.monotonic() + GLOBAL_TTL, prof) return prof async def _frontier_depths(db) -> tuple[list[dict], dict | None]: """Distribution des profondeurs + domaine le plus profond (scans lourds sur frontier_items ~3,7 M lignes → cache 1 h).""" hit = _global_cache.get("frontier_depths") if hit and hit[0] > time.monotonic(): return hit[1] # type: ignore[return-value] rows = await db.pool.fetch( "SELECT least(depth, 10)::int AS d, count(*)::int AS n FROM frontier_items GROUP BY 1 ORDER BY 1" ) bins = [{"label": ("10+" if r["d"] >= 10 else str(r["d"])), "value": r["n"]} for r in rows] deep = await db.pool.fetchrow( """ SELECT d.domain, f.depth FROM frontier_items f JOIN urls u ON u.id = f.url_id JOIN domains d ON d.id = u.domain_id ORDER BY f.depth DESC LIMIT 1 """ ) deepest = {"domain": deep["domain"], "depth": deep["depth"]} if deep else None val = (bins, deepest) _global_cache["frontier_depths"] = (time.monotonic() + DEEP_TTL, val) return val async def _tld_breakdown(db) -> list[dict]: """Pages indexées par TLD (table domains, petite — cache 10 min).""" hit = _global_cache.get("tld") if hit and hit[0] > time.monotonic(): return hit[1] # type: ignore[return-value] rows = await db.pool.fetch( """ SELECT CASE WHEN domain LIKE '%.qc.ca' THEN 'qc.ca' ELSE substring(domain from '[^.]+$') END AS tld, sum(page_count)::int AS n FROM domains WHERE page_count > 0 GROUP BY 1 ORDER BY n DESC LIMIT 9 """ ) items = [{"label": "." + r["tld"], "value": r["n"]} for r in rows if r["tld"]] _global_cache["tld"] = (time.monotonic() + GLOBAL_TTL, items) return items async def _crawl_hour_record(db) -> dict | None: """Heure record de crawl sur 30 jours (plage indexée — cache 10 min).""" hit = _global_cache.get("hour_record") if hit and hit[0] > time.monotonic(): return hit[1] # type: ignore[return-value] since = datetime.now(TZ) - timedelta(days=30) row = await db.pool.fetchrow( """ SELECT date_trunc('hour', fetched_at AT TIME ZONE 'America/Toronto') AS h, count(*)::int AS n FROM crawl_attempts WHERE fetched_at >= $1 GROUP BY 1 ORDER BY n DESC LIMIT 1 """, since, ) val = {"hour": row["h"], "n": row["n"]} if row else None _global_cache["hour_record"] = (time.monotonic() + GLOBAL_TTL, val) return val async def _purge_old_queries(db) -> None: """Rétention 90 jours des journaux de recherche — au plus 1×/jour.""" global _last_purge if time.monotonic() - _last_purge < 86400 and _last_purge: return _last_purge = time.monotonic() try: res = await db.pool.execute( "DELETE FROM search_queries WHERE created_at < now() - interval '90 days'" ) log.info("purge search_queries (90 j)", extra={"ctx": {"result": res}}) except Exception: log.exception("purge search_queries échouée") async def dashboard(db, period: str = "30j", from_s: str | None = None, to_s: str | None = None) -> dict: key = f"{period}|{from_s or ''}|{to_s or ''}" hit = _cache.get(key) if hit and hit[0] > time.monotonic(): return hit[1] async with _build_lock: hit = _cache.get(key) if hit and hit[0] > time.monotonic(): return hit[1] data = await _build(db, period, from_s, to_s) ttl = CACHE_TTL_TOUT if period == "tout" else CACHE_TTL _cache[key] = (time.monotonic() + ttl, data) return data async def _build(db, period: str, from_s: str | None, to_s: str | None) -> dict: now = datetime.now(TZ) ds_start = await _dataset_start(db) start, end, label, prev_start, prev_end = resolve_period(period, from_s, to_s, ds_start) bucket = "hour" if (end - start) <= timedelta(hours=36) else "day" unit_sql = "hour" if bucket == "hour" else "day" unit_label = "par heure" if bucket == "hour" else "par jour" hours = max((end - start).total_seconds() / 3600.0, 0.01) keys, klabels = _bucket_grid(start, end, bucket) await _purge_old_queries(db) # ---- Crawl : un passage borné, ventilé par (bucket, outcome) ---------- oc_rows = await db.pool.fetch( f""" SELECT date_trunc('{unit_sql}', fetched_at AT TIME ZONE 'America/Toronto') AS b, outcome, count(*)::int AS n, coalesce(sum(bytes), 0)::bigint AS bytes, max(bytes)::int AS bytes_max, sum(duration_ms)::bigint AS dur_sum, max(duration_ms)::int AS dur_max FROM crawl_attempts WHERE fetched_at >= $1 AND fetched_at < $2 GROUP BY 1, 2 """, start, end, ) oc_map: dict[tuple, int] = {} outcome_totals: dict[str, int] = {} attempts_cur = errors_cur = indexed_ok_cur = 0 bytes_sum = dur_sum = 0 bytes_max = dur_max = 0 for r in oc_rows: oc_map[(r["b"], r["outcome"])] = r["n"] outcome_totals[r["outcome"]] = outcome_totals.get(r["outcome"], 0) + r["n"] attempts_cur += r["n"] bytes_sum += r["bytes"] dur_sum += r["dur_sum"] or 0 bytes_max = max(bytes_max, r["bytes_max"] or 0) dur_max = max(dur_max, r["dur_max"] or 0) errors_cur = outcome_totals.get("error", 0) indexed_ok_cur = outcome_totals.get("indexed", 0) crawl_pts = [{"t": lab, "v": sum(oc_map.get((k, o), 0) for o in outcome_totals)} for k, lab in zip(keys, klabels)] err_pts = [{"t": lab, "v": oc_map.get((k, "error"), 0)} for k, lab in zip(keys, klabels)] idx_ok_pts = [{"t": lab, "v": oc_map.get((k, "indexed"), 0)} for k, lab in zip(keys, klabels)] # ---- Crawl : classes de statut HTTP par bucket (empilées) ------------- st_rows = await db.pool.fetch( f""" SELECT date_trunc('{unit_sql}', fetched_at AT TIME ZONE 'America/Toronto') AS b, CASE WHEN status_code BETWEEN 200 AND 299 THEN '2xx' WHEN status_code BETWEEN 300 AND 399 THEN '3xx' WHEN status_code BETWEEN 400 AND 499 THEN '4xx' WHEN status_code >= 500 THEN '5xx' ELSE 'Sans réponse' END AS cls, count(*)::int AS n FROM crawl_attempts WHERE fetched_at >= $1 AND fetched_at < $2 GROUP BY 1, 2 """, start, end, ) st_map = {(r["b"], r["cls"]): r["n"] for r in st_rows} st_totals = {c: sum(n for (b, cls), n in st_map.items() if cls == c) for c in HTTP_CLASSES} st_keys = [c for c in HTTP_CLASSES if st_totals.get(c)] stacked_pts = [{"t": lab, "values": [st_map.get((k, c), 0) for c in st_keys]} for k, lab in zip(keys, klabels)] # ---- Période précédente (bornée) --------------------------------------- prev_oc: dict[str, int] = {} attempts_prev = errors_prev = indexed_prev = searches_prev = bytes_prev = None if prev_start is not None: prows = await db.pool.fetch( """ SELECT outcome, count(*)::int AS n, coalesce(sum(bytes), 0)::bigint AS bytes FROM crawl_attempts WHERE fetched_at >= $1 AND fetched_at < $2 GROUP BY 1 """, prev_start, prev_end, ) prev_oc = {r["outcome"]: r["n"] for r in prows} attempts_prev = sum(prev_oc.values()) errors_prev = prev_oc.get("error", 0) bytes_prev = sum(r["bytes"] for r in prows) prev = await db.pool.fetchrow( """ SELECT (SELECT count(*)::int FROM documents WHERE first_indexed_at >= $1 AND first_indexed_at < $2) AS indexed_prev, (SELECT count(*)::int FROM search_queries WHERE created_at >= $1 AND created_at < $2) AS searches_prev """, prev_start, prev_end, ) indexed_prev = prev["indexed_prev"] searches_prev = prev["searches_prev"] # ---- Index / domaines / frontier (compteurs indexés) ------------------- cur = await db.pool.fetchrow( """ SELECT (SELECT count(*)::int FROM domains WHERE page_count > 0) AS domains_count, (SELECT count(*)::int FROM documents WHERE first_indexed_at >= $1 AND first_indexed_at < $2) AS indexed_cur, (SELECT count(*)::int FROM documents WHERE first_indexed_at < $1) AS docs_before, (SELECT count(*)::int FROM frontier_items WHERE status = 'pending') AS frontier_pending """, start, end, ) prof = await _doc_profile(db) docs_total = prof["total"] indexed_cur = cur["indexed_cur"] # ---- Recherches --------------------------------------------------------- sq = await db.pool.fetchrow( """ SELECT count(*)::int AS n, round(avg((NOT zero_result)::int) * 100, 1)::float AS found_rate, round(avg(took_ms))::int AS took_avg FROM search_queries WHERE created_at >= $1 AND created_at < $2 """, start, end, ) # ---- Séries ------------------------------------------------------------- per_bucket = await _bucket_counts(db, "documents", "first_indexed_at", start, end, bucket) running = cur["docs_before"] cum_points = [] for p in per_bucket: running += p["v"] cum_points.append({"t": p["t"], "v": running}) idx_compare = None if prev_start is not None and bucket == "day": idx_compare = await _bucket_counts( db, "documents", "first_indexed_at", prev_start, prev_end, bucket ) growth = round(100.0 * indexed_cur / cur["docs_before"], 1) if cur["docs_before"] else None vol_val, vol_unit = _mb(bytes_sum) kpis = [ {"id": "pages", "label": "Pages indexées (total)", "value": docs_total, "delta_pct": growth, "direction": "up", "spark": _spark(cum_points)}, {"id": "domaines", "label": "Sites du Groupe KA", "value": cur["domains_count"]}, {"id": "crawlees", "label": f"Pages crawlées — {label}", "value": attempts_cur, "delta_pct": _delta(attempts_cur, attempts_prev), "spark": _spark(crawl_pts)}, {"id": "rythme", "label": "Rythme de crawl", "value": round(attempts_cur / hours, 1), "unit": "pages/h", "delta_pct": _delta(attempts_cur, attempts_prev)}, {"id": "indexees", "label": f"Pages indexées — {label}", "value": indexed_cur, "delta_pct": _delta(indexed_cur, indexed_prev), "spark": _spark(per_bucket)}, {"id": "erreurs", "label": f"Erreurs de crawl — {label}", "value": errors_cur, "delta_pct": _delta(errors_cur, errors_prev), "spark": _spark(err_pts)}, {"id": "frontier", "label": "Frontier en attente", "value": cur["frontier_pending"]}, {"id": "recherches", "label": f"Recherches — {label}", "value": sq["n"], "delta_pct": _delta(sq["n"], searches_prev)}, ] if bytes_sum: kpis.append({"id": "volume", "label": f"Volume téléchargé — {label}", "value": vol_val, "unit": vol_unit, "delta_pct": _delta(bytes_sum, bytes_prev)}) if dur_sum and attempts_cur: kpis.append({"id": "duree", "label": "Durée moyenne d'une tentative", "value": int(round(dur_sum / attempts_cur)), "unit": "ms"}) # ---- Jauges ------------------------------------------------------------- gauges = [] if attempts_cur: gauges.append({"id": "succes", "label": "Crawl sans erreur", "value": round(100.0 * (attempts_cur - errors_cur) / attempts_cur, 1), "max": 100, "unit": "%", "help": "Part des tentatives de crawl qui ne finissent pas en erreur"}) gauges.append({"id": "convert", "label": "Crawl menant à l'index", "value": round(100.0 * indexed_ok_cur / attempts_cur, 1), "max": 100, "unit": "%", "help": "Part des tentatives dont l'issue est « indexée »"}) if docs_total: gauges.append({"id": "frais", "label": "Index frais (< 30 jours)", "value": round(100.0 * prof["fresh30"] / docs_total, 1), "max": 100, "unit": "%", "help": "Pages re-visitées par le crawler dans les 30 derniers jours"}) if sq["n"] and sq["found_rate"] is not None: gauges.append({"id": "taux", "label": "Recherches avec résultats", "value": sq["found_rate"], "max": 100, "unit": "%"}) series = [ {"id": "cumul", "title": "Pages indexées — cumul", "unit": "pages", "kind": "area", "points": cum_points}, {"id": "crawlees", "title": f"Pages crawlées {unit_label}", "unit": "pages", "kind": "bar", "points": crawl_pts}, {"id": "indexees", "title": f"Pages indexées {unit_label}", "unit": "pages", "kind": "line", "points": per_bucket, **({"compare": idx_compare} if idx_compare else {})}, {"id": "erreurs", "title": f"Erreurs de crawl {unit_label}", "unit": "erreurs", "kind": "line", "points": err_pts}, ] if sq["n"]: sq_bucket = await _bucket_counts(db, "search_queries", "created_at", start, end, bucket) series.append({"id": "recherches", "title": f"Recherches {unit_label}", "unit": "recherches", "kind": "bar", "points": sq_bucket}) multiseries = [{ "id": "crawlmix", "title": f"Crawl {unit_label} : tentatives, indexées, erreurs", "unit": "pages", "series": [ {"label": "Crawlées", "points": crawl_pts}, {"label": "Indexées", "points": idx_ok_pts}, {"label": "Erreurs", "points": err_pts}, ], }] stacked = [] if st_keys: stacked.append({"id": "http", "title": f"Crawl {unit_label} par classe de statut HTTP", "unit": "pages", "keys": st_keys, "points": stacked_pts}) # ---- Répartitions ------------------------------------------------------- frontier = await db.pool.fetch( "SELECT status, count(*)::int AS n FROM frontier_items GROUP BY status ORDER BY n DESC" ) top_domains = await db.pool.fetch( """ SELECT domain, page_count, url_count, round(quebec_score::numeric, 2)::float AS qs, last_crawled_at FROM domains WHERE page_count > 0 ORDER BY page_count DESC LIMIT 100 """ ) tld_items = await _tld_breakdown(db) breakdowns = [ {"id": "outcomes", "title": f"Résultats de crawl — {label}", "kind": "donut", "items": [{"label": OUTCOME_LABELS.get(o, o), "value": n, **({"delta_pct": _delta(n, prev_oc.get(o))} if prev_oc else {})} for o, n in sorted(outcome_totals.items(), key=lambda kv: -kv[1])]}, {"id": "http", "title": f"Crawl par classe de statut HTTP — {label}", "kind": "bars", "items": [{"label": c, "value": st_totals[c]} for c in st_keys]}, {"id": "topdom", "title": "Pages indexées par site du Groupe KA (tout l'index)", "kind": "bars", "items": [{"label": r["domain"], "value": r["page_count"]} for r in top_domains[:14]]}, {"id": "frontier", "title": "File frontier par statut", "kind": "donut", "items": [{"label": FRONTIER_LABELS.get(r["status"], r["status"]), "value": r["n"]} for r in frontier]}, ] if tld_items: breakdowns.insert(2, {"id": "tld", "title": "Pages indexées par TLD", "kind": "donut", "items": tld_items}) if prof["langs"]: breakdowns.append({"id": "langues", "title": "Langue des pages indexées", "kind": "donut", "items": prof["langs"][:8]}) # ---- Distributions ------------------------------------------------------ distributions = [] size_rows = await db.pool.fetch( """ SELECT width_bucket(bytes, 0, 512000, 16) AS wb, count(*)::int AS n FROM crawl_attempts WHERE fetched_at >= $1 AND fetched_at < $2 AND bytes IS NOT NULL AND bytes > 0 GROUP BY 1 ORDER BY 1 """, start, end, ) if size_rows: bins = [] got = {r["wb"]: r["n"] for r in size_rows} for i in range(1, 17): bins.append({"label": f"{(i - 1) * 32}-{i * 32} Ko", "value": got.get(i, 0)}) over = got.get(17, 0) if over: bins.append({"label": "> 512 Ko", "value": over}) distributions.append({"id": "taille", "title": f"Distribution de la taille des pages — {label}", "unit": "pages", "bins": bins}) dur_rows = await db.pool.fetch( """ SELECT width_bucket(duration_ms, 0, 10000, 10) AS wb, count(*)::int AS n FROM crawl_attempts WHERE fetched_at >= $1 AND fetched_at < $2 AND duration_ms IS NOT NULL GROUP BY 1 ORDER BY 1 """, start, end, ) if dur_rows: got = {r["wb"]: r["n"] for r in dur_rows} bins = [{"label": f"{i - 1}-{i} s", "value": got.get(i, 0)} for i in range(1, 11)] over = got.get(11, 0) if over: bins.append({"label": "> 10 s", "value": over}) distributions.append({"id": "duree", "title": f"Distribution de la durée des requêtes — {label}", "unit": "pages", "bins": bins}) depth_bins, deepest = await _frontier_depths(db) if depth_bins: distributions.append({"id": "profondeur", "title": "Profondeur des URLs de la frontier", "unit": "URLs", "bins": depth_bins}) # ---- Heatmap horaire 7×24 (bornée à 56 jours max) ----------------------- hh_start = max(start, end - timedelta(days=56)) hh_rows = await db.pool.fetch( """ SELECT (extract(isodow FROM fetched_at AT TIME ZONE 'America/Toronto')::int - 1) AS dow, extract(hour FROM fetched_at AT TIME ZONE 'America/Toronto')::int AS hr, count(*)::int AS n FROM crawl_attempts WHERE fetched_at >= $1 AND fetched_at < $2 GROUP BY 1, 2 """, hh_start, end, ) hourly = None if hh_rows: hourly = {"title": "Pages crawlées par jour × heure", "cells": [{"dow": r["dow"], "hour": r["hr"], "value": r["n"]} for r in hh_rows]} # ---- Tableaux ----------------------------------------------------------- tables = [ {"id": "topdom", "title": "Sites du Groupe KA (tout l'index)", "columns": ["Site", "Pages indexées", "URLs connues", "Dernier crawl"], "rows": [[r["domain"], r["page_count"], r["url_count"], r["last_crawled_at"].astimezone(TZ).strftime("%Y-%m-%d %H:%M") if r["last_crawled_at"] else "—"] for r in top_domains]}, ] top_q = [] if sq["n"]: top_q = await db.pool.fetch( """ SELECT lower(query) AS q, count(*)::int AS n, round(avg(results_total))::int AS avg_res, max(created_at) AS last_at FROM search_queries WHERE created_at >= $1 AND created_at < $2 GROUP BY 1 ORDER BY n DESC, last_at DESC LIMIT 50 """, start, end, ) tables.append({ "id": "topq", "title": f"Top requêtes de recherche — {label}", "columns": ["Requête", "Recherches", "Résultats moy.", "Dernière fois"], "rows": [[r["q"], r["n"], r["avg_res"], r["last_at"].astimezone(TZ).strftime("%Y-%m-%d %H:%M")] for r in top_q], }) zero_q = await db.pool.fetch( """ SELECT lower(query) AS q, count(*)::int AS n, max(created_at) AS last_at FROM search_queries WHERE zero_result AND created_at >= $1 AND created_at < $2 GROUP BY 1 ORDER BY n DESC, last_at DESC LIMIT 30 """, start, end, ) if zero_q: tables.append({ "id": "zeroq", "title": f"Requêtes sans résultat — {label}", "columns": ["Requête", "Tentatives", "Dernière fois"], "rows": [[r["q"], r["n"], r["last_at"].astimezone(TZ).strftime("%Y-%m-%d %H:%M")] for r in zero_q], }) err_dom = await db.pool.fetch( """ SELECT d.domain, count(*)::int AS n, max(a.fetched_at) AS last_at FROM crawl_attempts a JOIN urls u ON u.id = a.url_id JOIN domains d ON d.id = u.domain_id WHERE a.outcome = 'error' AND a.fetched_at >= $1 AND a.fetched_at < $2 GROUP BY 1 ORDER BY n DESC LIMIT 30 """, start, end, ) if err_dom: tables.append({ "id": "errdom", "title": f"Erreurs de crawl par domaine — {label}", "columns": ["Domaine", "Erreurs", "Dernière erreur"], "rows": [[r["domain"], r["n"], r["last_at"].astimezone(TZ).strftime("%Y-%m-%d %H:%M")] for r in err_dom], }) err_types = await db.pool.fetch( """ SELECT coalesce(error_code, 'inconnu') AS code, count(*)::int AS n, max(fetched_at) AS last_at FROM crawl_attempts WHERE outcome = 'error' AND fetched_at >= $1 AND fetched_at < $2 GROUP BY 1 ORDER BY n DESC """, start, end, ) if err_types: tables.append({ "id": "errtypes", "title": f"Erreurs de crawl par type — {label}", "columns": ["Type d'erreur", "Occurrences", "Dernière occurrence"], "rows": [[r["code"], r["n"], r["last_at"].astimezone(TZ).strftime("%Y-%m-%d %H:%M")] for r in err_types], }) last_errors = await db.pool.fetch( """ SELECT a.fetched_at, coalesce(a.error_code, 'inconnu') AS code, u.url FROM crawl_attempts a JOIN urls u ON u.id = a.url_id WHERE a.outcome = 'error' AND a.fetched_at >= $1 AND a.fetched_at < $2 ORDER BY a.fetched_at DESC LIMIT 25 """, start, end, ) if last_errors: tables.append({ "id": "lasterr", "title": "Dernières erreurs de crawl", "columns": ["Quand", "Type", "URL"], "rows": [[r["fetched_at"].astimezone(TZ).strftime("%Y-%m-%d %H:%M"), r["code"], r["url"][:120]] for r in last_errors], }) # ---- Heatmap calendrier & records (réels, jamais inventés) -------------- cells = await _heatmap_cells(db) records = [] if cells: best = max(cells, key=lambda c: c["value"]) records.append({"label": "Jour record d'indexation", "value": f"{_fr_int(best['value'])} pages", "date": best["date"]}) hr_rec = await _crawl_hour_record(db) if hr_rec: records.append({"label": "Heure record de crawl (30 derniers jours)", "value": f"{_fr_int(hr_rec['n'])} pages", "date": hr_rec["hour"].strftime("%Y-%m-%d %Hh")}) if top_domains: records.append({"label": "Site le plus fourni", "value": f"{top_domains[0]['domain']} — {_fr_int(top_domains[0]['page_count'])} pages"}) if deepest: records.append({"label": "Site exploré le plus profondément", "value": f"{deepest['domain']} (profondeur {deepest['depth']})"}) if bytes_max: bm_val, bm_unit = _mb(bytes_max) records.append({"label": f"Page la plus lourde téléchargée — {label}", "value": f"{str(bm_val).replace('.', ',')} {bm_unit}"}) if dur_max: records.append({"label": f"Plus longue tentative de crawl — {label}", "value": f"{str(round(dur_max / 1000, 1)).replace('.', ',')} s"}) if prof["avg_qs"] is not None: records.append({"label": "Score Québec moyen des pages indexées", "value": str(prof["avg_qs"]).replace(".", ",")}) if top_q: records.append({"label": f"Requête la plus fréquente — {label}", "value": f"« {top_q[0]['q'][:40]} » — {_fr_int(top_q[0]['n'])} fois"}) if sq["n"] and sq["took_avg"] is not None: records.append({"label": f"Temps de réponse moyen des recherches — {label}", "value": f"{_fr_int(sq['took_avg'])} ms"}) if sq["n"]: best_sq = await db.pool.fetchrow( """ SELECT (created_at AT TIME ZONE 'America/Toronto')::date AS d, count(*)::int AS n FROM search_queries GROUP BY 1 ORDER BY n DESC LIMIT 1 """ ) if best_sq: records.append({"label": "Jour record de recherches", "value": f"{_fr_int(best_sq['n'])} recherches", "date": best_sq["d"].isoformat()}) out = { "updated": now.isoformat(), "period": {"from": start.astimezone(TZ).date().isoformat(), "to": end.astimezone(TZ).date().isoformat(), "label": label}, "kpis": kpis, "series": series, "multiseries": multiseries, "breakdowns": breakdowns, "heatmap": {"title": "Indexations par jour (26 dernières semaines)", "cells": cells}, "tables": tables, "records": records, } if gauges: out["gauges"] = gauges if stacked: out["stacked"] = stacked if distributions: out["distributions"] = distributions if hourly: out["hourly"] = hourly return out