404 lines
15 KiB
Python
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")
|