Files
ahad-quant/performance_monitor.py
2026-06-25 14:00:20 +03:00

404 lines
15 KiB
Python

"""
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")