""" AHAD QUANT — Performance Monitor (Apprentissage Continu V7) ============================================================== Surveille la performance en temps réel et détecte quand le modèle dégrade. Métriques surveillées (sliding window 50 trades) : - win_rate → alerte si < 45% - avg_pnl → alerte si < -0.002 - max_drawdown → alerte si > 8% - confidence_gap → alerte si confiance élevée mais pertes (overconfidence) - timeout_rate → alerte si > 70% (modèle trop incertain) Niveaux d'alerte : OK : tout va bien WARNING : win_rate < 50% pendant 3 jours → log + notification DANGER : win_rate < 45% → augmente MIN_CONFIDENCE automatiquement EMERGENCY : win_rate < 35% OU drawdown > 10% → pause trading Drift detection (Page-Hinkley) : - Détecte un changement de distribution sur le PnL ou le win_rate - Si dérive détectée → log + recommande retrain complet Usage : from experience_buffer import get_experience_buffer from performance_monitor import PerformanceMonitor monitor = PerformanceMonitor(buffer) status = monitor.check() # "OK" / "WARNING" / "DANGER" / "EMERGENCY" report = monitor.get_report() monitor.log_daily_report() """ import json import logging import os import time import threading from datetime import datetime, timezone from typing import Dict, List, Optional import numpy as np import config from experience_buffer import ExperienceBuffer from features import NUM_FEATURES log = logging.getLogger("PerformanceMonitor") # ── Constantes ──────────────────────────────────────────────────────────────── WINDOW_SIZE = 50 # trades pour sliding window WIN_RATE_WARNING = 0.50 # seuil warning WIN_RATE_DANGER = 0.45 # seuil danger WIN_RATE_EMERGENCY = 0.35 # seuil urgence (pause trading) DRAWDOWN_EMERGENCY = 0.10 # 10% drawdown = urgence TIMEOUT_RATE_WARNING = 0.70 # 70% timeout = modèle trop incertain AVG_PNL_WARNING = -0.002 # PnL moyen négatif = warning MONITOR_INTERVAL_SEC = 3600 # check toutes les heures WARNING_DAYS_THRESHOLD = 3 # jours consécutifs sous seuil → WARNING envoyé REPORT_FILE = "performance_report.json" class HealthStatus: OK = "OK" WARNING = "WARNING" DANGER = "DANGER" EMERGENCY = "EMERGENCY" class PerformanceMonitor: """ Surveillance temps réel des métriques de performance du bot. Détecte le drift et déclenche des alertes ou la pause du trading. """ def __init__( self, buffer: ExperienceBuffer, report_file: str = REPORT_FILE, notify_fn=None, ): self.buffer = buffer self.report_file = report_file self.notify_fn = notify_fn self._lock = threading.RLock() self._last_status = HealthStatus.OK self._warning_since: Optional[float] = None # timestamp premier WARNING self._ph_state = _PageHinkleyState() # détecteur de drift # ── Check principal ────────────────────────────────────────────────────── def check(self) -> str: """ Évalue la santé du système sur les WINDOW_SIZE derniers trades. Retourne : "OK" / "WARNING" / "DANGER" / "EMERGENCY" """ if not getattr(config, "MONITOR_ENABLED", True): return HealthStatus.OK recent = self.buffer.get_last_n(WINDOW_SIZE) if len(recent) < 10: return HealthStatus.OK # pas assez de données stats = self.buffer.stats(window=WINDOW_SIZE) status = self._evaluate(stats) # Mise à jour du drift detector if recent: last_pnl = recent[-1].get("pnl", 0.0) drift = self._ph_state.update(last_pnl) if drift: log.warning("[MONITOR] 🚨 DRIFT DÉTECTÉ (Page-Hinkley) — distribution PnL changée") self._notify( "⚠️ *AHAD QUANT DRIFT DÉTECTÉ*\n" "Distribution du PnL a changé significativement.\n" "→ Recommande un retrain complet (local ou cloud au choix)." ) # Gestion de l'état WARNING persistant if status == HealthStatus.WARNING: if self._warning_since is None: self._warning_since = time.time() elif time.time() - self._warning_since >= WARNING_DAYS_THRESHOLD * 86400: log.warning(f"[MONITOR] WARNING persistant depuis {WARNING_DAYS_THRESHOLD} jours") self._notify( f"⚠️ *AHAD QUANT WARNING persistant*\n" f"win_rate < 50% depuis {WARNING_DAYS_THRESHOLD}+ jours.\n" f"win_rate actuel : {stats['win_rate']:.1%}\n" f"→ Vérifier les conditions de marché." ) else: self._warning_since = None # reset si plus en WARNING # Actions selon le niveau if status == HealthStatus.DANGER: self._handle_danger(stats) elif status == HealthStatus.EMERGENCY: self._handle_emergency(stats) self._last_status = status return status def _evaluate(self, stats: Dict) -> str: """Détermine le niveau de santé selon les métriques.""" win_rate = stats.get("win_rate", 1.0) max_dd = stats.get("max_drawdown", 0.0) timeout_r = stats.get("timeout_rate", 0.0) avg_pnl = stats.get("avg_pnl", 0.0) # Urgence : critères stricts if win_rate < WIN_RATE_EMERGENCY or max_dd > DRAWDOWN_EMERGENCY: return HealthStatus.EMERGENCY # Danger : performance dégradée if win_rate < WIN_RATE_DANGER or avg_pnl < AVG_PNL_WARNING: return HealthStatus.DANGER # Warning : signaux précoces if win_rate < WIN_RATE_WARNING or timeout_r > TIMEOUT_RATE_WARNING: return HealthStatus.WARNING return HealthStatus.OK def _handle_danger(self, stats: Dict) -> None: """DANGER : augmenter MIN_CONFIDENCE pour être plus sélectif.""" current_conf = getattr(config, "MIN_CONFIDENCE", 0.72) new_conf = min(current_conf + 0.02, 0.88) config.MIN_CONFIDENCE = new_conf msg = ( f"🔴 *AHAD QUANT DANGER*\n" f"win_rate={stats['win_rate']:.1%} | avg_pnl={stats['avg_pnl']:.4f}\n" f"MIN_CONFIDENCE ajusté : {current_conf:.3f} → {new_conf:.3f}" ) log.warning(msg.replace("*","").replace("`","")) self._notify(msg) def _handle_emergency(self, stats: Dict) -> None: """EMERGENCY : pause du trading si activé.""" if not getattr(config, "EMERGENCY_PAUSE_ENABLED", True): return msg = ( f"🚨 *AHAD QUANT EMERGENCY*\n" f"win_rate={stats['win_rate']:.1%} | " f"max_dd={stats['max_drawdown']:.1%}\n" f"→ PAUSE trading activée.\n" f"→ Relancer un retrain complet (local ou cloud au choix)." ) log.critical(msg.replace("*","").replace("`","")) self._notify(msg) # Signaler la demande de pause via config config.PAPER_MODE = True # Basculer en paper mode comme filet de sécurité log.critical("[MONITOR] Bot basculé en PAPER MODE d'urgence") def should_emergency_retrain(self) -> bool: """True si un retrain complet d'urgence est recommandé.""" stats = self.buffer.stats(window=WINDOW_SIZE) return ( stats.get("win_rate", 1.0) < WIN_RATE_EMERGENCY or stats.get("max_drawdown", 0.0) > DRAWDOWN_EMERGENCY ) # ── Rapport ────────────────────────────────────────────────────────────── def get_report(self) -> Dict: """Retourne un rapport complet des métriques actuelles.""" stats_50 = self.buffer.stats(window=WINDOW_SIZE) stats_all = self.buffer.stats() recent = self.buffer.get_last_n(WINDOW_SIZE) # Confidence gap : trades avec haute confiance mais résultat LOSS high_conf_losses = [ e for e in recent if e.get("confidence", 0) > 0.80 and e.get("outcome") == "LOSS" ] confidence_gap = len(high_conf_losses) / max(len(recent), 1) return { "timestamp": datetime.now(timezone.utc).isoformat(), "status": self._last_status, "window": min(len(recent), WINDOW_SIZE), "metrics_50": stats_50, "metrics_all": stats_all, "confidence_gap": round(confidence_gap, 4), "drift_detected": self._ph_state.drift_detected, "min_confidence": getattr(config, "MIN_CONFIDENCE", 0.72), "thresholds": { "win_rate_warning": WIN_RATE_WARNING, "win_rate_danger": WIN_RATE_DANGER, "win_rate_emergency": WIN_RATE_EMERGENCY, "drawdown_emergency": DRAWDOWN_EMERGENCY, "timeout_warning": TIMEOUT_RATE_WARNING, }, } def log_daily_report(self) -> None: """Sauvegarde le rapport quotidien dans performance_report.json.""" report = self.get_report() try: tmp = self.report_file + ".tmp" with open(tmp, "w", encoding="utf-8") as f: json.dump(report, f, indent=2, ensure_ascii=False) os.replace(tmp, self.report_file) s = report["metrics_50"] log.info( f"[MONITOR] Rapport quotidien | status={report['status']} | " f"win_rate={s['win_rate']:.1%} | avg_pnl={s['avg_pnl']:.4f} | " f"dd={s['max_drawdown']:.1%} | trades={s['total_trades']}" ) except Exception as e: log.error(f"[MONITOR] Erreur sauvegarde rapport : {e}") @property def last_status(self) -> str: return self._last_status def _notify(self, msg: str) -> None: if self.notify_fn: try: self.notify_fn(msg) except Exception as e: log.debug(f"Erreur notification : {e}") # ── Page-Hinkley Drift Detector ─────────────────────────────────────────────── class _PageHinkleyState: """ Implémentation simple du test de Page-Hinkley pour détecter un changement de distribution (drift) sur une série de PnL. Déclenche si la somme cumulée dépasse un seuil λ (lambda_). """ def __init__(self, delta: float = 0.005, lambda_: float = 0.15, alpha: float = 0.9999): self.delta = delta # sensibilité au changement self.lambda_ = lambda_ # seuil de déclenchement self.alpha = alpha # facteur d'oubli self._sum = 0.0 self._min_sum = 0.0 self._n = 0 self._mean = 0.0 self.drift_detected = False def update(self, value: float) -> bool: """ Mise à jour avec une nouvelle observation. Retourne True si un drift est détecté. """ self._n += 1 # Mise à jour de la moyenne en ligne (avec oubli) self._mean = self.alpha * self._mean + (1 - self.alpha) * value # Somme cumulée avec biais δ self._sum += (self._mean - value - self.delta) self._min_sum = min(self._min_sum, self._sum) # Test de Page-Hinkley if self._n > 30 and (self._sum - self._min_sum) > self.lambda_: self.drift_detected = True # Reset après détection self._sum = 0.0 self._min_sum = 0.0 return True self.drift_detected = False return False # ── Background Monitor Thread ───────────────────────────────────────────────── class MonitorThread: """ Thread en arrière-plan qui appelle monitor.check() toutes les heures et monitor.log_daily_report() une fois par jour. """ def __init__(self, monitor: PerformanceMonitor): self.monitor = monitor self._stop = threading.Event() self._thread = None self._last_daily = 0.0 def start(self) -> None: if not getattr(config, "MONITOR_ENABLED", True): return self._thread = threading.Thread( target=self._loop, daemon=True, name="PerformanceMonitor" ) self._thread.start() log.info("[MONITOR] Thread démarré (check toutes les heures)") def stop(self) -> None: self._stop.set() def _loop(self) -> None: while not self._stop.is_set(): try: status = self.monitor.check() log.debug(f"[MONITOR] Status : {status}") # Rapport quotidien toutes les 24h if time.time() - self._last_daily >= 86400: self.monitor.log_daily_report() self._last_daily = time.time() except Exception as e: log.error(f"[MONITOR] Erreur loop : {e}") self._stop.wait(timeout=MONITOR_INTERVAL_SEC) # ── CLI de test ─────────────────────────────────────────────────────────────── if __name__ == "__main__": import logging logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(name)s] %(message)s") from experience_buffer import ExperienceBuffer import numpy as np buf = ExperienceBuffer(max_size=200, buffer_file="test_monitor_buffer.json") # Simuler 60 trades avec mauvaise performance for i in range(60): outcome = "LOSS" if i % 3 != 0 else "WIN" # 33% win rate → DANGER buf.add({ "pair": "EURUSD", "signal": "LONG", "confidence": 0.75, "features": list(np.random.randn(NUM_FEATURES).astype(float)), "regime": "VOLATILE", "entry_price": 1.0850, "exit_price": 1.0820, "pnl": 0.02 if outcome == "WIN" else -0.015, "outcome": outcome, "hold_candles": 4, }) monitor = PerformanceMonitor(buf, report_file="test_performance_report.json") status = monitor.check() print(f"\nStatus : {status}") report = monitor.get_report() print(f"\nRapport :") print(f" win_rate = {report['metrics_50']['win_rate']:.1%}") print(f" avg_pnl = {report['metrics_50']['avg_pnl']:.4f}") print(f" status = {report['status']}") print(f" drift = {report['drift_detected']}") monitor.log_daily_report() print(f"\nRapport sauvegardé : test_performance_report.json") # Nettoyage for f in ["test_monitor_buffer.json", "test_performance_report.json"]: if os.path.exists(f): os.remove(f) print("\n✅ PerformanceMonitor — test OK")