# ----------------------------------------------------------------------------- # Groupe KA — kaid.py : client KA ID v2 (personnalisation) pour les satellites. # SOURCE CANONIQUE : ka-ui.git/kaid/kaid.py — copié dans le paquet backend de # chaque app (louka/, jobka/, sortika/, …) par sync-kaid.sh. Ne pas diverger : # corriger ICI puis redistribuer. # # Rôle : relier l'app au feature store du hub (groupe-ka.com) — # · track() journal d'interactions (serveur, fil d'exécution dédié) # · fetch_prefs() profil de préférences appris (cache 90 s, fail-open) # · rerank() reclassement personnalisé APRÈS la pertinence de base # · build_router() routes /api/kaid/* (événements client, masquage, # recherches sauvegardées) # # Contrat s2s (identique à hubfav/hubprofile) : HMAC-SHA256 du secret SSO # partagé — sig = HMAC(KA_SSO_SECRET, f"{CLIENT_ID}.{ka_id}.{ts}"). # Config .env : KA_SSO_SECRET (déjà présent), KA_HUB_URL (optionnel). # # Principes : la connexion n'est JAMAIS requise ; sans profil ou à la moindre # erreur réseau → classement de base inchangé (fail-open). La personnalisation # ne remplace pas la pertinence : elle reclasse (blend) et n'écrase jamais # l'intention de la session (les dimensions explicitement filtrées par la # requête courante sont ignorées dans le score). # ----------------------------------------------------------------------------- from __future__ import annotations import hashlib import hmac import json import os import threading import time import requests from fastapi import APIRouter, HTTPException, Request from pydantic import BaseModel KA_HUB_URL = os.environ.get("KA_HUB_URL", "https://www.groupe-ka.com").rstrip("/") CLIENT_ID = os.environ.get("KA_CLIENT_ID", "") # fixé par init() dans web.py PREFS_TTL = 90 # secondes de cache du profil TIMEOUT = 5 # secondes par appel hub LOCATION_DIMS = {"city", "region", "sector", "quartier", "ville", "location", "neighborhood"} LANGUAGE_DIMS = {"language", "langue"} PRICE_DIMS = {"price", "rent", "salary", "salary_year", "price_min"} # Événements acceptés depuis le navigateur (le reste vient du serveur). CLIENT_EVENT_TYPES = { "click", "impression", "detail_dwell", "scroll_depth", "return_visit", "share", "compare", "external_click", "map_open", "map_marker_click", "alert_open", } _prefs_cache: dict[str, tuple[float, dict | None]] = {} _seen_searches: dict[str, float] = {} # anti-doublon des recherches (120 s) _lock = threading.Lock() def init(client_id: str) -> None: """À appeler une fois au démarrage de l'app (web.py).""" global CLIENT_ID CLIENT_ID = client_id def _sig(ka_id: str, ts: int) -> str | None: secret = os.environ.get("KA_SSO_SECRET") if not secret or not CLIENT_ID: return None return hmac.new(secret.encode(), f"{CLIENT_ID}.{ka_id}.{ts}".encode(), hashlib.sha256).hexdigest() def _signed_params(ka_id: str) -> dict | None: ts = int(time.time()) sig = _sig(ka_id, ts) if not sig: return None return {"client_id": CLIENT_ID, "ka_id": ka_id, "ts": str(ts), "sig": sig} def _ka_id_of(user) -> str | None: """Extrait un ka_id exploitable d'un dict utilisateur (ou None).""" if not user: return None ka = (user.get("ka_id") or "").strip() if isinstance(user, dict) else "" return ka if ka.startswith("ka-") else None # ---------------------------------------------------------------- événements def _post_events(ka_id: str, events: list[dict]) -> None: p = _signed_params(ka_id) if not p: return try: requests.post(f"{KA_HUB_URL}/api/sso/events", timeout=TIMEOUT, json={**p, "events": events}) except Exception: pass # best-effort : jamais bloquant, jamais fatal def track(user, etype: str, *, entity_type: str | None = None, entity_id: str | None = None, query: str | None = None, filters: dict | None = None, position: int | None = None, features: dict | None = None, dwell_ms: int | None = None, session_id: str | None = None) -> None: """Journalise un événement au hub (fil dédié, zéro latence ajoutée). No-op si l'utilisateur n'est pas connecté via KA ID.""" ka_id = _ka_id_of(user) if not ka_id: return if etype == "search": # anti-rafale : la même recherche (mêmes filtres) < 120 s n'est # journalisée qu'une fois — une SPA relance l'API à chaque frappe. key = ka_id + "|" + hashlib.sha1( json.dumps([query, filters], sort_keys=True, default=str).encode() ).hexdigest() now = time.time() with _lock: if now - _seen_searches.get(key, 0) < 120: return _seen_searches[key] = now if len(_seen_searches) > 2000: cutoff = now - 300 for k in [k for k, t in _seen_searches.items() if t < cutoff]: del _seen_searches[k] ev: dict = {"type": etype} if entity_type: ev["entity_type"] = entity_type if entity_id: ev["entity_id"] = str(entity_id) if query: ev["query"] = str(query)[:200] if filters: ev["filters"] = filters if position is not None: ev["position"] = int(position) if features: ev["features"] = features if dwell_ms is not None: ev["dwell_ms"] = int(dwell_ms) if session_id: ev["session_id"] = str(session_id)[:60] threading.Thread(target=_post_events, args=(ka_id, [ev]), daemon=True).start() # ------------------------------------------------------------------ profil def fetch_prefs(ka_id: str | None) -> dict | None: """Profil de personnalisation du membre (cache 90 s). None si non connecté, non configuré ou hub injoignable — l'appelant retombe alors sur le classement de base.""" if not ka_id or not str(ka_id).startswith("ka-"): return None now = time.time() with _lock: hit = _prefs_cache.get(ka_id) if hit and now - hit[0] < PREFS_TTL: return hit[1] data: dict | None = None p = _signed_params(ka_id) if p: try: r = requests.get(f"{KA_HUB_URL}/api/sso/prefs", params=p, timeout=TIMEOUT) if r.status_code == 200: data = r.json() except Exception: data = None with _lock: _prefs_cache[ka_id] = (now, data) if len(_prefs_cache) > 500: for k in list(_prefs_cache)[:100]: del _prefs_cache[k] return data def invalidate_prefs(ka_id: str | None) -> None: if not ka_id: return with _lock: _prefs_cache.pop(ka_id, None) # ---------------------------------------------------------------- reranking def _norm(v) -> str: return str(v).strip().lower() def personal_score(feats: dict, app_profile: dict, global_profile: dict, active_dims: set[str]) -> tuple[float | None, list[str]]: """Score personnel [0,1] d'une annonce, ou None si le profil ne couvre aucune de ses caractéristiques. `active_dims` = dimensions explicitement filtrées par la requête courante (intention de session > long terme).""" dims = app_profile.get("dims") or {} ranges = app_profile.get("ranges") or {} gl = (global_profile or {}).get("location") or {} num = 0.0 den = 0.0 reasons: list[str] = [] for dim, val in (feats or {}).items(): if val is None or dim in active_dims: continue if isinstance(val, bool): val = str(val) if isinstance(val, (int, float)): r = ranges.get(dim) if r and r.get("n", 0) >= 5: p25, p75 = r["p25"], r["p75"] iqr = max(p75 - p25, abs(r.get("p50", 0)) * 0.1, 1.0) if p25 <= val <= p75: aff = 1.0 elif p25 - 1.5 * iqr <= val <= p75 + 1.5 * iqr: aff = 0.3 else: aff = -0.4 # poids réduit : une plage numérique seule (prix…) ne doit # jamais suffire à personnaliser (0.6 < seuil den 0.8) — # sinon tout item au « bon prix » score 1.0 et noie les # correspondances réelles (ville, marque, type). w = 0.6 num += w * aff den += w if aff == 1.0: reasons.append("MATCH_PRICE_RANGE" if dim in PRICE_DIMS else f"MATCH_{dim.upper()}_RANGE") continue vals = val if isinstance(val, (list, tuple)) else [val] vals = [_norm(v) for v in vals if v not in (None, "")] if not vals: continue d = dims.get(dim) if d: vv = d.get("values") or {} affs = [vv[v] for v in vals if v in vv] if affs: aff = max(affs) w = float(d.get("conf") or 0.5) num += w * aff den += w if aff >= 0.6: reasons.append("MATCH_LOCATION" if dim in LOCATION_DIMS else f"MATCH_{dim.upper()}") if dim in LOCATION_DIMS: gv = gl.get("values") or {} affs = [gv[v] for v in vals if v in gv] if affs and max(affs) > 0: w = 0.6 * float(gl.get("conf") or 0.3) num += w * max(affs) den += w if max(affs) >= 0.6 and "MATCH_LOCATION" not in reasons: reasons.append("MATCH_LOCATION") if dim in LANGUAGE_DIMS: glang = (global_profile or {}).get("language") or {} gv = glang.get("values") or {} affs = [gv[v] for v in vals if v in gv] if affs and max(affs) > 0: w = 0.4 * float(glang.get("conf") or 0.3) num += w * max(affs) den += w if den < 0.8: return None, [] score = (num / den + 1.0) / 2.0 return max(0.0, min(1.0, score)), reasons[:4] def rerank(items: list, user, *, features_of, uid_of=None, active_dims: set[str] | None = None, blend: float = 0.35, badge: float = 0.62, max_considered: int = 300, reco_key: str = "ka_reco") -> tuple[list, bool]: """Reclassement personnalisé APRÈS la pertinence de base. · items : liste (dicts) déjà triée par la pertinence de base · features_of : item -> dict de caractéristiques {dim: valeur} · uid_of : item -> identifiant canonique (défaut : item["uid"]) · active_dims : dimensions filtrées par la requête (ignorées du score) Retourne (items, personnalisé?). Les annonces masquées (« Pas pour moi ») sont retirées. Annote item[reco_key] = {score, reasons} quand le score personnel est net (badge « Recommandé pour vous » — parcimonieux).""" if uid_of is None: uid_of = lambda it: (it.get("uid") if isinstance(it, dict) else None) ka_id = _ka_id_of(user) if not ka_id or not items: return items, False prefs = fetch_prefs(ka_id) if not prefs: return items, False hidden = set(prefs.get("hidden") or []) if hidden: items = [it for it in items if str(uid_of(it)) not in hidden] if not prefs.get("personalization"): return items, False profile = prefs.get("profile") or {} app_p = profile.get("app") if not app_p or not items: return items, False head = items[:max_considered] tail = items[max_considered:] n = len(head) active = active_dims or set() # signaux collaboratifs du hub : co-favoris (item-item) et # recommandations du modèle de matrix factorization (ALS, batch quotidien) similar = {str(s) for s in (app_p.get("similar") or [])} mf = {str(s) for s in (app_p.get("mf") or [])} scored = [] badged = 0 for i, it in enumerate(head): base = 1.0 - i / max(n, 1) try: p, reasons = personal_score(features_of(it) or {}, app_p, profile.get("global") or {}, active) except Exception: p, reasons = None, [] uid = str(uid_of(it)) if similar and uid in similar: p = min(1.0, (p if p is not None else 0.55) + 0.25) reasons = (["SIMILAR_USERS"] + reasons)[:4] elif mf and uid in mf: p = min(1.0, (p if p is not None else 0.55) + 0.25) reasons = (["COLLABORATIVE_MODEL"] + reasons)[:4] if p is None: final = (1.0 - blend) * base + blend * 0.5 else: final = (1.0 - blend) * base + blend * p if p >= badge and reasons and badged < max(2, n // 8) \ and isinstance(it, dict): it[reco_key] = {"score": round(p, 2), "reasons": reasons} badged += 1 scored.append((final, i, it)) scored.sort(key=lambda t: (-t[0], t[1])) # stable : départage par rang return [it for _, _, it in scored] + tail, True # ------------------------------------------------------- proxys hub (s2s) def _hub_post(ka_id: str, path: str, payload: dict) -> dict: p = _signed_params(ka_id) if not p: raise HTTPException(503, "KA_SSO_SECRET manquant (voir .env)") try: r = requests.post(f"{KA_HUB_URL}{path}", timeout=TIMEOUT, json={**p, **payload}) return r.json() if r.status_code == 200 else {"error": r.status_code} except Exception: raise HTTPException(502, "hub KA injoignable") def _hub_get(ka_id: str, path: str) -> dict: p = _signed_params(ka_id) if not p: raise HTTPException(503, "KA_SSO_SECRET manquant (voir .env)") try: r = requests.get(f"{KA_HUB_URL}{path}", params=p, timeout=TIMEOUT) return r.json() if r.status_code == 200 else {"error": r.status_code} except Exception: raise HTTPException(502, "hub KA injoignable") # ------------------------------------------------------------------ routeur class _EventsIn(BaseModel): events: list[dict] class _HideIn(BaseModel): item_id: str on: bool = True features: dict | None = None class _SearchIn(BaseModel): action: str = "add" # add | remove | alert | touch id: int | None = None label: str | None = None query: str | None = None filters: dict | None = None location: str | None = None url: str | None = None alert: bool = False frequency: str | None = None def build_router(get_user) -> APIRouter: """Routes /api/kaid/* de l'app. `get_user(request)` = current_user de l'app (dict avec ka_id, ou None).""" router = APIRouter(prefix="/api/kaid") def _require_ka(request: Request) -> tuple[dict, str]: user = get_user(request) ka_id = _ka_id_of(user) if not ka_id: raise HTTPException(401, "connexion KA ID requise") return user, ka_id @router.get("/status") def status(request: Request): user = get_user(request) ka_id = _ka_id_of(user) if not ka_id: return {"connected": False} prefs = fetch_prefs(ka_id) return { "connected": True, "personalization": bool(prefs and prefs.get("personalization")), "monka_url": f"{KA_HUB_URL}/mon-ka", } @router.post("/events") def client_events(request: Request, body: _EventsIn): user = get_user(request) ka_id = _ka_id_of(user) if not ka_id: return {"ok": True, "stored": 0} events = [] for e in body.events[:20]: if e.get("type") in CLIENT_EVENT_TYPES: events.append({k: e[k] for k in ("type", "entity_type", "entity_id", "query", "filters", "position", "features", "dwell_ms", "session_id") if k in e}) if events: threading.Thread(target=_post_events, args=(ka_id, events), daemon=True).start() return {"ok": True, "stored": len(events)} @router.post("/hide") def hide(request: Request, body: _HideIn): _, ka_id = _require_ka(request) out = _hub_post(ka_id, "/api/sso/hide", { "item_id": body.item_id, "on": body.on, "features": body.features, }) invalidate_prefs(ka_id) return out @router.get("/saved-searches") def saved_list(request: Request): _, ka_id = _require_ka(request) return _hub_get(ka_id, "/api/sso/saved-searches") @router.post("/saved-searches") def saved_post(request: Request, body: _SearchIn): _, ka_id = _require_ka(request) search: dict = {k: v for k, v in { "id": body.id, "label": body.label, "query": body.query, "filters": body.filters, "location": body.location, "url": body.url, "alert": body.alert, "frequency": body.frequency, }.items() if v is not None} out = _hub_post(ka_id, "/api/sso/saved-searches", {"action": body.action, "search": search}) invalidate_prefs(ka_id) return out return router