|
1 |
+# ----------------------------------------------------------------------------- |
|
2 |
+# Groupe KA — kaid.py : client KA ID v2 (personnalisation) pour les satellites. |
|
3 |
+# SOURCE CANONIQUE : ka-ui.git/kaid/kaid.py — copié dans le paquet backend de |
|
4 |
+# chaque app (louka/, jobka/, sortika/, …) par sync-kaid.sh. Ne pas diverger : |
|
5 |
+# corriger ICI puis redistribuer. |
|
6 |
+# |
|
7 |
+# Rôle : relier l'app au feature store du hub (groupe-ka.com) — |
|
8 |
+# · track() journal d'interactions (serveur, fil d'exécution dédié) |
|
9 |
+# · fetch_prefs() profil de préférences appris (cache 90 s, fail-open) |
|
10 |
+# · rerank() reclassement personnalisé APRÈS la pertinence de base |
|
11 |
+# · build_router() routes /api/kaid/* (événements client, masquage, |
|
12 |
+# recherches sauvegardées) |
|
13 |
+# |
|
14 |
+# Contrat s2s (identique à hubfav/hubprofile) : HMAC-SHA256 du secret SSO |
|
15 |
+# partagé — sig = HMAC(KA_SSO_SECRET, f"{CLIENT_ID}.{ka_id}.{ts}"). |
|
16 |
+# Config .env : KA_SSO_SECRET (déjà présent), KA_HUB_URL (optionnel). |
|
17 |
+# |
|
18 |
+# Principes : la connexion n'est JAMAIS requise ; sans profil ou à la moindre |
|
19 |
+# erreur réseau → classement de base inchangé (fail-open). La personnalisation |
|
20 |
+# ne remplace pas la pertinence : elle reclasse (blend) et n'écrase jamais |
|
21 |
+# l'intention de la session (les dimensions explicitement filtrées par la |
|
22 |
+# requête courante sont ignorées dans le score). |
|
23 |
+# ----------------------------------------------------------------------------- |
|
24 |
+from __future__ import annotations |
|
25 |
+ |
|
26 |
+import hashlib |
|
27 |
+import hmac |
|
28 |
+import json |
|
29 |
+import os |
|
30 |
+import threading |
|
31 |
+import time |
|
32 |
+ |
|
33 |
+import requests |
|
34 |
+from fastapi import APIRouter, HTTPException, Request |
|
35 |
+from pydantic import BaseModel |
|
36 |
+ |
|
37 |
+KA_HUB_URL = os.environ.get("KA_HUB_URL", "https://www.groupe-ka.com").rstrip("/") |
|
38 |
+CLIENT_ID = os.environ.get("KA_CLIENT_ID", "") # fixé par init() dans web.py |
|
39 |
+ |
|
40 |
+PREFS_TTL = 90 # secondes de cache du profil |
|
41 |
+TIMEOUT = 5 # secondes par appel hub |
|
42 |
+LOCATION_DIMS = {"city", "region", "sector", "quartier", "ville", |
|
43 |
+ "location", "neighborhood"} |
|
44 |
+PRICE_DIMS = {"price", "rent", "salary", "salary_year", "price_min"} |
|
45 |
+ |
|
46 |
+# Événements acceptés depuis le navigateur (le reste vient du serveur). |
|
47 |
+CLIENT_EVENT_TYPES = { |
|
48 |
+ "click", "impression", "detail_dwell", "scroll_depth", "return_visit", |
|
49 |
+ "share", "compare", "external_click", "map_open", "map_marker_click", |
|
50 |
+ "alert_open", |
|
51 |
+} |
|
52 |
+ |
|
53 |
+_prefs_cache: dict[str, tuple[float, dict | None]] = {} |
|
54 |
+_seen_searches: dict[str, float] = {} # anti-doublon des recherches (120 s) |
|
55 |
+_lock = threading.Lock() |
|
56 |
+ |
|
57 |
+ |
|
58 |
+def init(client_id: str) -> None: |
|
59 |
+ """À appeler une fois au démarrage de l'app (web.py).""" |
|
60 |
+ global CLIENT_ID |
|
61 |
+ CLIENT_ID = client_id |
|
62 |
+ |
|
63 |
+ |
|
64 |
+def _sig(ka_id: str, ts: int) -> str | None: |
|
65 |
+ secret = os.environ.get("KA_SSO_SECRET") |
|
66 |
+ if not secret or not CLIENT_ID: |
|
67 |
+ return None |
|
68 |
+ return hmac.new(secret.encode(), |
|
69 |
+ f"{CLIENT_ID}.{ka_id}.{ts}".encode(), |
|
70 |
+ hashlib.sha256).hexdigest() |
|
71 |
+ |
|
72 |
+ |
|
73 |
+def _signed_params(ka_id: str) -> dict | None: |
|
74 |
+ ts = int(time.time()) |
|
75 |
+ sig = _sig(ka_id, ts) |
|
76 |
+ if not sig: |
|
77 |
+ return None |
|
78 |
+ return {"client_id": CLIENT_ID, "ka_id": ka_id, "ts": str(ts), "sig": sig} |
|
79 |
+ |
|
80 |
+ |
|
81 |
+def _ka_id_of(user) -> str | None: |
|
82 |
+ """Extrait un ka_id exploitable d'un dict utilisateur (ou None).""" |
|
83 |
+ if not user: |
|
84 |
+ return None |
|
85 |
+ ka = (user.get("ka_id") or "").strip() if isinstance(user, dict) else "" |
|
86 |
+ return ka if ka.startswith("ka-") else None |
|
87 |
+ |
|
88 |
+ |
|
89 |
+# ---------------------------------------------------------------- événements |
|
90 |
+ |
|
91 |
+def _post_events(ka_id: str, events: list[dict]) -> None: |
|
92 |
+ p = _signed_params(ka_id) |
|
93 |
+ if not p: |
|
94 |
+ return |
|
95 |
+ try: |
|
96 |
+ requests.post(f"{KA_HUB_URL}/api/sso/events", timeout=TIMEOUT, |
|
97 |
+ json={**p, "events": events}) |
|
98 |
+ except Exception: |
|
99 |
+ pass # best-effort : jamais bloquant, jamais fatal |
|
100 |
+ |
|
101 |
+ |
|
102 |
+def track(user, etype: str, *, entity_type: str | None = None, |
|
103 |
+ entity_id: str | None = None, query: str | None = None, |
|
104 |
+ filters: dict | None = None, position: int | None = None, |
|
105 |
+ features: dict | None = None, dwell_ms: int | None = None, |
|
106 |
+ session_id: str | None = None) -> None: |
|
107 |
+ """Journalise un événement au hub (fil dédié, zéro latence ajoutée). |
|
108 |
+ No-op si l'utilisateur n'est pas connecté via KA ID.""" |
|
109 |
+ ka_id = _ka_id_of(user) |
|
110 |
+ if not ka_id: |
|
111 |
+ return |
|
112 |
+ if etype == "search": |
|
113 |
+ # anti-rafale : la même recherche (mêmes filtres) < 120 s n'est |
|
114 |
+ # journalisée qu'une fois — une SPA relance l'API à chaque frappe. |
|
115 |
+ key = ka_id + "|" + hashlib.sha1( |
|
116 |
+ json.dumps([query, filters], sort_keys=True, default=str).encode() |
|
117 |
+ ).hexdigest() |
|
118 |
+ now = time.time() |
|
119 |
+ with _lock: |
|
120 |
+ if now - _seen_searches.get(key, 0) < 120: |
|
121 |
+ return |
|
122 |
+ _seen_searches[key] = now |
|
123 |
+ if len(_seen_searches) > 2000: |
|
124 |
+ cutoff = now - 300 |
|
125 |
+ for k in [k for k, t in _seen_searches.items() if t < cutoff]: |
|
126 |
+ del _seen_searches[k] |
|
127 |
+ ev: dict = {"type": etype} |
|
128 |
+ if entity_type: ev["entity_type"] = entity_type |
|
129 |
+ if entity_id: ev["entity_id"] = str(entity_id) |
|
130 |
+ if query: ev["query"] = str(query)[:200] |
|
131 |
+ if filters: ev["filters"] = filters |
|
132 |
+ if position is not None: ev["position"] = int(position) |
|
133 |
+ if features: ev["features"] = features |
|
134 |
+ if dwell_ms is not None: ev["dwell_ms"] = int(dwell_ms) |
|
135 |
+ if session_id: ev["session_id"] = str(session_id)[:60] |
|
136 |
+ threading.Thread(target=_post_events, args=(ka_id, [ev]), daemon=True).start() |
|
137 |
+ |
|
138 |
+ |
|
139 |
+# ------------------------------------------------------------------ profil |
|
140 |
+ |
|
141 |
+def fetch_prefs(ka_id: str | None) -> dict | None: |
|
142 |
+ """Profil de personnalisation du membre (cache 90 s). None si non |
|
143 |
+ connecté, non configuré ou hub injoignable — l'appelant retombe alors |
|
144 |
+ sur le classement de base.""" |
|
145 |
+ if not ka_id or not str(ka_id).startswith("ka-"): |
|
146 |
+ return None |
|
147 |
+ now = time.time() |
|
148 |
+ with _lock: |
|
149 |
+ hit = _prefs_cache.get(ka_id) |
|
150 |
+ if hit and now - hit[0] < PREFS_TTL: |
|
151 |
+ return hit[1] |
|
152 |
+ data: dict | None = None |
|
153 |
+ p = _signed_params(ka_id) |
|
154 |
+ if p: |
|
155 |
+ try: |
|
156 |
+ r = requests.get(f"{KA_HUB_URL}/api/sso/prefs", params=p, |
|
157 |
+ timeout=TIMEOUT) |
|
158 |
+ if r.status_code == 200: |
|
159 |
+ data = r.json() |
|
160 |
+ except Exception: |
|
161 |
+ data = None |
|
162 |
+ with _lock: |
|
163 |
+ _prefs_cache[ka_id] = (now, data) |
|
164 |
+ if len(_prefs_cache) > 500: |
|
165 |
+ for k in list(_prefs_cache)[:100]: |
|
166 |
+ del _prefs_cache[k] |
|
167 |
+ return data |
|
168 |
+ |
|
169 |
+ |
|
170 |
+def invalidate_prefs(ka_id: str | None) -> None: |
|
171 |
+ if not ka_id: |
|
172 |
+ return |
|
173 |
+ with _lock: |
|
174 |
+ _prefs_cache.pop(ka_id, None) |
|
175 |
+ |
|
176 |
+ |
|
177 |
+# ---------------------------------------------------------------- reranking |
|
178 |
+ |
|
179 |
+def _norm(v) -> str: |
|
180 |
+ return str(v).strip().lower() |
|
181 |
+ |
|
182 |
+ |
|
183 |
+def personal_score(feats: dict, app_profile: dict, global_profile: dict, |
|
184 |
+ active_dims: set[str]) -> tuple[float | None, list[str]]: |
|
185 |
+ """Score personnel [0,1] d'une annonce, ou None si le profil ne couvre |
|
186 |
+ aucune de ses caractéristiques. `active_dims` = dimensions explicitement |
|
187 |
+ filtrées par la requête courante (intention de session > long terme).""" |
|
188 |
+ dims = app_profile.get("dims") or {} |
|
189 |
+ ranges = app_profile.get("ranges") or {} |
|
190 |
+ gl = (global_profile or {}).get("location") or {} |
|
191 |
+ num = 0.0 |
|
192 |
+ den = 0.0 |
|
193 |
+ reasons: list[str] = [] |
|
194 |
+ for dim, val in (feats or {}).items(): |
|
195 |
+ if val is None or dim in active_dims: |
|
196 |
+ continue |
|
197 |
+ if isinstance(val, bool): |
|
198 |
+ val = str(val) |
|
199 |
+ if isinstance(val, (int, float)): |
|
200 |
+ r = ranges.get(dim) |
|
201 |
+ if r and r.get("n", 0) >= 5: |
|
202 |
+ p25, p75 = r["p25"], r["p75"] |
|
203 |
+ iqr = max(p75 - p25, abs(r.get("p50", 0)) * 0.1, 1.0) |
|
204 |
+ if p25 <= val <= p75: |
|
205 |
+ aff = 1.0 |
|
206 |
+ elif p25 - 1.5 * iqr <= val <= p75 + 1.5 * iqr: |
|
207 |
+ aff = 0.3 |
|
208 |
+ else: |
|
209 |
+ aff = -0.4 |
|
210 |
+ num += aff |
|
211 |
+ den += 1.0 |
|
212 |
+ if aff == 1.0: |
|
213 |
+ reasons.append("MATCH_PRICE_RANGE" if dim in PRICE_DIMS |
|
214 |
+ else f"MATCH_{dim.upper()}_RANGE") |
|
215 |
+ continue |
|
216 |
+ vals = val if isinstance(val, (list, tuple)) else [val] |
|
217 |
+ vals = [_norm(v) for v in vals if v not in (None, "")] |
|
218 |
+ if not vals: |
|
219 |
+ continue |
|
220 |
+ d = dims.get(dim) |
|
221 |
+ if d: |
|
222 |
+ vv = d.get("values") or {} |
|
223 |
+ affs = [vv[v] for v in vals if v in vv] |
|
224 |
+ if affs: |
|
225 |
+ aff = max(affs) |
|
226 |
+ w = float(d.get("conf") or 0.5) |
|
227 |
+ num += w * aff |
|
228 |
+ den += w |
|
229 |
+ if aff >= 0.6: |
|
230 |
+ reasons.append("MATCH_LOCATION" if dim in LOCATION_DIMS |
|
231 |
+ else f"MATCH_{dim.upper()}") |
|
232 |
+ if dim in LOCATION_DIMS: |
|
233 |
+ gv = gl.get("values") or {} |
|
234 |
+ affs = [gv[v] for v in vals if v in gv] |
|
235 |
+ if affs and max(affs) > 0: |
|
236 |
+ w = 0.6 * float(gl.get("conf") or 0.3) |
|
237 |
+ num += w * max(affs) |
|
238 |
+ den += w |
|
239 |
+ if max(affs) >= 0.6 and "MATCH_LOCATION" not in reasons: |
|
240 |
+ reasons.append("MATCH_LOCATION") |
|
241 |
+ if den < 0.8: |
|
242 |
+ return None, [] |
|
243 |
+ score = (num / den + 1.0) / 2.0 |
|
244 |
+ return max(0.0, min(1.0, score)), reasons[:4] |
|
245 |
+ |
|
246 |
+ |
|
247 |
+def rerank(items: list, user, *, features_of, uid_of=None, |
|
248 |
+ active_dims: set[str] | None = None, blend: float = 0.35, |
|
249 |
+ badge: float = 0.62, max_considered: int = 300, |
|
250 |
+ reco_key: str = "ka_reco") -> tuple[list, bool]: |
|
251 |
+ """Reclassement personnalisé APRÈS la pertinence de base. |
|
252 |
+ · items : liste (dicts) déjà triée par la pertinence de base |
|
253 |
+ · features_of : item -> dict de caractéristiques {dim: valeur} |
|
254 |
+ · uid_of : item -> identifiant canonique (défaut : item["uid"]) |
|
255 |
+ · active_dims : dimensions filtrées par la requête (ignorées du score) |
|
256 |
+ Retourne (items, personnalisé?). Les annonces masquées (« Pas pour moi ») |
|
257 |
+ sont retirées. Annote item[reco_key] = {score, reasons} quand le score |
|
258 |
+ personnel est net (badge « Recommandé pour vous » — parcimonieux).""" |
|
259 |
+ if uid_of is None: |
|
260 |
+ uid_of = lambda it: (it.get("uid") if isinstance(it, dict) else None) |
|
261 |
+ ka_id = _ka_id_of(user) |
|
262 |
+ if not ka_id or not items: |
|
263 |
+ return items, False |
|
264 |
+ prefs = fetch_prefs(ka_id) |
|
265 |
+ if not prefs: |
|
266 |
+ return items, False |
|
267 |
+ hidden = set(prefs.get("hidden") or []) |
|
268 |
+ if hidden: |
|
269 |
+ items = [it for it in items if str(uid_of(it)) not in hidden] |
|
270 |
+ if not prefs.get("personalization"): |
|
271 |
+ return items, False |
|
272 |
+ profile = prefs.get("profile") or {} |
|
273 |
+ app_p = profile.get("app") |
|
274 |
+ if not app_p or not items: |
|
275 |
+ return items, False |
|
276 |
+ |
|
277 |
+ head = items[:max_considered] |
|
278 |
+ tail = items[max_considered:] |
|
279 |
+ n = len(head) |
|
280 |
+ active = active_dims or set() |
|
281 |
+ scored = [] |
|
282 |
+ badged = 0 |
|
283 |
+ for i, it in enumerate(head): |
|
284 |
+ base = 1.0 - i / max(n, 1) |
|
285 |
+ try: |
|
286 |
+ p, reasons = personal_score(features_of(it) or {}, app_p, |
|
287 |
+ profile.get("global") or {}, active) |
|
288 |
+ except Exception: |
|
289 |
+ p, reasons = None, [] |
|
290 |
+ if p is None: |
|
291 |
+ final = (1.0 - blend) * base + blend * 0.5 |
|
292 |
+ else: |
|
293 |
+ final = (1.0 - blend) * base + blend * p |
|
294 |
+ if p >= badge and reasons and badged < max(2, n // 8) \ |
|
295 |
+ and isinstance(it, dict): |
|
296 |
+ it[reco_key] = {"score": round(p, 2), "reasons": reasons} |
|
297 |
+ badged += 1 |
|
298 |
+ scored.append((final, i, it)) |
|
299 |
+ scored.sort(key=lambda t: (-t[0], t[1])) # stable : départage par rang |
|
300 |
+ return [it for _, _, it in scored] + tail, True |
|
301 |
+ |
|
302 |
+ |
|
303 |
+# ------------------------------------------------------- proxys hub (s2s) |
|
304 |
+ |
|
305 |
+def _hub_post(ka_id: str, path: str, payload: dict) -> dict: |
|
306 |
+ p = _signed_params(ka_id) |
|
307 |
+ if not p: |
|
308 |
+ raise HTTPException(503, "KA_SSO_SECRET manquant (voir .env)") |
|
309 |
+ try: |
|
310 |
+ r = requests.post(f"{KA_HUB_URL}{path}", timeout=TIMEOUT, |
|
311 |
+ json={**p, **payload}) |
|
312 |
+ return r.json() if r.status_code == 200 else {"error": r.status_code} |
|
313 |
+ except Exception: |
|
314 |
+ raise HTTPException(502, "hub KA injoignable") |
|
315 |
+ |
|
316 |
+ |
|
317 |
+def _hub_get(ka_id: str, path: str) -> dict: |
|
318 |
+ p = _signed_params(ka_id) |
|
319 |
+ if not p: |
|
320 |
+ raise HTTPException(503, "KA_SSO_SECRET manquant (voir .env)") |
|
321 |
+ try: |
|
322 |
+ r = requests.get(f"{KA_HUB_URL}{path}", params=p, timeout=TIMEOUT) |
|
323 |
+ return r.json() if r.status_code == 200 else {"error": r.status_code} |
|
324 |
+ except Exception: |
|
325 |
+ raise HTTPException(502, "hub KA injoignable") |
|
326 |
+ |
|
327 |
+ |
|
328 |
+# ------------------------------------------------------------------ routeur |
|
329 |
+ |
|
330 |
+class _EventsIn(BaseModel): |
|
331 |
+ events: list[dict] |
|
332 |
+ |
|
333 |
+ |
|
334 |
+class _HideIn(BaseModel): |
|
335 |
+ item_id: str |
|
336 |
+ on: bool = True |
|
337 |
+ features: dict | None = None |
|
338 |
+ |
|
339 |
+ |
|
340 |
+class _SearchIn(BaseModel): |
|
341 |
+ action: str = "add" # add | remove | alert | touch |
|
342 |
+ id: int | None = None |
|
343 |
+ label: str | None = None |
|
344 |
+ query: str | None = None |
|
345 |
+ filters: dict | None = None |
|
346 |
+ location: str | None = None |
|
347 |
+ url: str | None = None |
|
348 |
+ alert: bool = False |
|
349 |
+ frequency: str | None = None |
|
350 |
+ |
|
351 |
+ |
|
352 |
+def build_router(get_user) -> APIRouter: |
|
353 |
+ """Routes /api/kaid/* de l'app. `get_user(request)` = current_user de |
|
354 |
+ l'app (dict avec ka_id, ou None).""" |
|
355 |
+ router = APIRouter(prefix="/api/kaid") |
|
356 |
+ |
|
357 |
+ def _require_ka(request: Request) -> tuple[dict, str]: |
|
358 |
+ user = get_user(request) |
|
359 |
+ ka_id = _ka_id_of(user) |
|
360 |
+ if not ka_id: |
|
361 |
+ raise HTTPException(401, "connexion KA ID requise") |
|
362 |
+ return user, ka_id |
|
363 |
+ |
|
364 |
+ @router.get("/status") |
|
365 |
+ def status(request: Request): |
|
366 |
+ user = get_user(request) |
|
367 |
+ ka_id = _ka_id_of(user) |
|
368 |
+ if not ka_id: |
|
369 |
+ return {"connected": False} |
|
370 |
+ prefs = fetch_prefs(ka_id) |
|
371 |
+ return { |
|
372 |
+ "connected": True, |
|
373 |
+ "personalization": bool(prefs and prefs.get("personalization")), |
|
374 |
+ "monka_url": f"{KA_HUB_URL}/mon-ka", |
|
375 |
+ } |
|
376 |
+ |
|
377 |
+ @router.post("/events") |
|
378 |
+ def client_events(request: Request, body: _EventsIn): |
|
379 |
+ user = get_user(request) |
|
380 |
+ ka_id = _ka_id_of(user) |
|
381 |
+ if not ka_id: |
|
382 |
+ return {"ok": True, "stored": 0} |
|
383 |
+ events = [] |
|
384 |
+ for e in body.events[:20]: |
|
385 |
+ if e.get("type") in CLIENT_EVENT_TYPES: |
|
386 |
+ events.append({k: e[k] for k in |
|
387 |
+ ("type", "entity_type", "entity_id", "query", |
|
388 |
+ "filters", "position", "features", "dwell_ms", |
|
389 |
+ "session_id") if k in e}) |
|
390 |
+ if events: |
|
391 |
+ threading.Thread(target=_post_events, args=(ka_id, events), |
|
392 |
+ daemon=True).start() |
|
393 |
+ return {"ok": True, "stored": len(events)} |
|
394 |
+ |
|
395 |
+ @router.post("/hide") |
|
396 |
+ def hide(request: Request, body: _HideIn): |
|
397 |
+ _, ka_id = _require_ka(request) |
|
398 |
+ out = _hub_post(ka_id, "/api/sso/hide", { |
|
399 |
+ "item_id": body.item_id, "on": body.on, |
|
400 |
+ "features": body.features, |
|
401 |
+ }) |
|
402 |
+ invalidate_prefs(ka_id) |
|
403 |
+ return out |
|
404 |
+ |
|
405 |
+ @router.get("/saved-searches") |
|
406 |
+ def saved_list(request: Request): |
|
407 |
+ _, ka_id = _require_ka(request) |
|
408 |
+ return _hub_get(ka_id, "/api/sso/saved-searches") |
|
409 |
+ |
|
410 |
+ @router.post("/saved-searches") |
|
411 |
+ def saved_post(request: Request, body: _SearchIn): |
|
412 |
+ _, ka_id = _require_ka(request) |
|
413 |
+ search: dict = {k: v for k, v in { |
|
414 |
+ "id": body.id, "label": body.label, "query": body.query, |
|
415 |
+ "filters": body.filters, "location": body.location, |
|
416 |
+ "url": body.url, "alert": body.alert, |
|
417 |
+ "frequency": body.frequency, |
|
418 |
+ }.items() if v is not None} |
|
419 |
+ out = _hub_post(ka_id, "/api/sso/saved-searches", |
|
420 |
+ {"action": body.action, "search": search}) |
|
421 |
+ invalidate_prefs(ka_id) |
|
422 |
+ return out |
|
423 |
+ |
|
424 |
+ return router |