279 lines
11 KiB
Python
279 lines
11 KiB
Python
"""
|
|
AHAD QUANT — Ensemble Core (Source Unique de Vérité ML)
|
|
============================================================
|
|
UNE seule implémentation de « comment transformer un ensemble ML brut
|
|
(LightGBM + XGBoost + RandomForest [+ TFT + TransformerGRU]) en une
|
|
probabilité calibrée ». Utilisée IDENTIQUEMENT par :
|
|
|
|
ahad_quant.py → signal live (paper / live / MT5)
|
|
backtest.py → validation historique
|
|
rl_env.py → observation de l'agent RL pendant l'ENTRAÎNEMENT
|
|
rl_agent.py → observation de l'agent RL en INFÉRENCE live
|
|
|
|
AVANT ce module, la même logique était ré-écrite à la main 3 fois, avec de
|
|
vraies divergences silencieuses :
|
|
|
|
• rl_env.py ignorait totalement TFT/TransformerGRU même quand has_dl=True,
|
|
et retombait en silence sur une simple moyenne des 3 modèles tabulaires
|
|
si le nombre de colonnes ne correspondait pas à scaler.n_features_in_
|
|
(le RL "voyait" donc un ML différent de celui réellement utilisé en
|
|
live dès que les modèles séquentiels étaient activés).
|
|
• ahad_quant.py pouvait lever une ValueError du StandardScaler (et donc
|
|
planter toute la boucle principale) si has_dl=True mais l'historique
|
|
disponible pour une paire était trop court pour les modèles séquentiels.
|
|
• backtest.py n'avait AUCUNE gestion de has_dl — le backtest validait un
|
|
modèle différent de celui réellement déployé.
|
|
|
|
Avec ce module, les 4 points d'entrée appellent EXACTEMENT le même code,
|
|
dans le même ordre de stacking que train.py (lgbm, xgb, rf, tft, tgru) :
|
|
le ML vu par le RL pendant l'entraînement == le ML vu en live == le ML
|
|
validé en backtest. Un seul cerveau, pas trois.
|
|
"""
|
|
|
|
import os
|
|
import logging
|
|
import pickle
|
|
import numpy as np
|
|
|
|
log = logging.getLogger("EnsembleCore")
|
|
|
|
# Valeur de repli si prepare_sequences n'est pas importable (ne devrait
|
|
# jamais arriver en pratique — prepare_sequences.py ne dépend que de numpy).
|
|
_DL_SEQ_LEN_FALLBACK = 168
|
|
|
|
|
|
# ─── Chargement ──────────────────────────────────────────────────────────────
|
|
|
|
def load_ensemble(path: str):
|
|
"""Charge un ensemble .pkl. Retourne None si absent/erreur (fail-safe —
|
|
ne lève jamais d'exception, pour que tous les appelants puissent
|
|
continuer en mode dégradé identique)."""
|
|
if not path or not os.path.exists(path):
|
|
log.warning(f"[ENS] Ensemble introuvable : {path}")
|
|
return None
|
|
try:
|
|
with open(path, "rb") as f:
|
|
ens = pickle.load(f)
|
|
acc = ens.get("ens_acc")
|
|
acc_str = f"{acc:.2%}" if isinstance(acc, (int, float)) else "?"
|
|
log.info(f"[ENS] Chargé ({path}) — has_dl={ens.get('has_dl', False)}, acc={acc_str}")
|
|
return ens
|
|
except Exception as e:
|
|
log.error(f"[ENS] Erreur chargement ensemble {path} : {e}")
|
|
return None
|
|
|
|
|
|
def ensemble_mtime(path: str) -> float:
|
|
"""mtime du fichier ensemble, 0.0 si absent (utilisé pour le hot-reload)."""
|
|
try:
|
|
return os.path.getmtime(path)
|
|
except OSError:
|
|
return 0.0
|
|
|
|
|
|
# ─── Modèles séquentiels (DL) ─────────────────────────────────────────────────
|
|
|
|
def has_dl_models(ensemble: dict) -> bool:
|
|
"""True seulement si l'ensemble déclare has_dl ET que les deux modèles
|
|
séquentiels sont réellement présents (évite tout état incohérent)."""
|
|
return bool(
|
|
ensemble is not None
|
|
and ensemble.get("has_dl")
|
|
and ensemble.get("tft") is not None
|
|
and ensemble.get("tgru") is not None
|
|
)
|
|
|
|
|
|
def dl_seq_len() -> int:
|
|
try:
|
|
from prepare_sequences import SEQ_LEN
|
|
return SEQ_LEN
|
|
except Exception:
|
|
return _DL_SEQ_LEN_FALLBACK
|
|
|
|
|
|
def build_sequence_windows(X: np.ndarray, seq_len: int):
|
|
"""
|
|
Construit TOUTES les fenêtres glissantes valides de longueur seq_len
|
|
pour l'inférence (contrairement à prepare_sequences.build_sequences, qui
|
|
est pour l'ENTRAÎNEMENT et exclut volontairement les dernières lignes
|
|
pour éviter le data-leakage des labels — ici il n'y a pas de label, on
|
|
veut une prédiction pour CHAQUE ligne disposant d'assez d'historique,
|
|
y compris la toute dernière).
|
|
|
|
Retourne (start_idx, X_seq) :
|
|
start_idx : index dans X de la 1ère ligne couverte (= seq_len - 1)
|
|
X_seq : (M, seq_len, F) avec M = len(X) - seq_len + 1
|
|
(None, None) si len(X) < seq_len (pas assez d'historique).
|
|
"""
|
|
N = len(X)
|
|
if N < seq_len:
|
|
return None, None
|
|
M = N - seq_len + 1
|
|
F = X.shape[1]
|
|
X_seq = np.empty((M, seq_len, F), dtype=np.float32)
|
|
for i in range(M):
|
|
X_seq[i] = X[i : i + seq_len]
|
|
return seq_len - 1, X_seq
|
|
|
|
|
|
def _predict_dl_batch(ensemble: dict, raw_feats: np.ndarray):
|
|
"""
|
|
Calcule (tft_probs, tgru_probs) ALIGNÉS sur raw_feats (même longueur N).
|
|
Les lignes sans historique suffisant reçoivent 0.5 (neutre) — jamais de
|
|
NaN, jamais d'exception qui remonte à l'appelant.
|
|
Retourne (None, None) si les modèles DL ne sont pas dispo / pas chargeables
|
|
(torch absent, erreur d'inférence, etc.) — l'appelant doit alors
|
|
simplement ne pas inclure TFT/TGRU dans le stacking, exactement comme un
|
|
sous-modèle manquant.
|
|
"""
|
|
if not has_dl_models(ensemble):
|
|
return None, None
|
|
try:
|
|
import torch
|
|
from prepare_sequences import transform_sequences
|
|
from tft_model import predict_tft_proba
|
|
from transformer_gru_model import predict_tgru_proba
|
|
except ImportError:
|
|
log.debug("[ENS] torch / tft_model / transformer_gru_model indisponibles — DL ignoré")
|
|
return None, None
|
|
|
|
seq_len = dl_seq_len()
|
|
feats = np.nan_to_num(raw_feats.astype(np.float32), nan=0.0, posinf=0.0, neginf=0.0)
|
|
start_idx, X_seq = build_sequence_windows(feats, seq_len)
|
|
|
|
N = len(raw_feats)
|
|
if X_seq is None:
|
|
# Pas assez d'historique pour CE pair → neutre partout, pas d'erreur.
|
|
return (np.full(N, 0.5, dtype=np.float32),
|
|
np.full(N, 0.5, dtype=np.float32))
|
|
|
|
dl_scaler = ensemble.get("dl_scaler")
|
|
if dl_scaler is not None:
|
|
try:
|
|
X_seq = transform_sequences(dl_scaler, X_seq)
|
|
except Exception as e:
|
|
log.warning(f"[ENS] Erreur normalisation séquence DL : {e}")
|
|
return None, None
|
|
|
|
try:
|
|
X_seq_t = torch.FloatTensor(X_seq)
|
|
tft_p = predict_tft_proba(ensemble["tft"], X_seq_t)
|
|
tgru_p = predict_tgru_proba(ensemble["tgru"], X_seq_t)
|
|
except Exception as e:
|
|
log.warning(f"[ENS] Erreur inférence DL (TFT/TGRU) : {e}")
|
|
return None, None
|
|
|
|
tft_full = np.full(N, 0.5, dtype=np.float32)
|
|
tgru_full = np.full(N, 0.5, dtype=np.float32)
|
|
tft_full[start_idx : start_idx + len(tft_p)] = tft_p
|
|
tgru_full[start_idx : start_idx + len(tgru_p)] = tgru_p
|
|
return tft_full, tgru_full
|
|
|
|
|
|
# ─── Fonction canonique : ensemble → (proba, confidence) ──────────────────────
|
|
|
|
def predict_ensemble_batch(ensemble: dict, raw_feats: np.ndarray, use_dl: bool = True) -> np.ndarray:
|
|
"""
|
|
UNIQUE implémentation : ensemble brut → (N, 2) [proba_long, confidence].
|
|
|
|
raw_feats : (N, NUM_FEATURES) — features BRUTES non normalisées (les
|
|
modèles tabulaires LightGBM/XGBoost/RF n'ont pas besoin de
|
|
scaling ; seules les features envoyées au réseau PPO le sont,
|
|
via un scaler totalement séparé — voir rl_agent.py).
|
|
|
|
Stacking dans le MÊME ordre que train.py::build_ensemble (déterminant
|
|
pour la cohérence du méta-modèle) : lgbm, xgb, rf, [tft, tgru].
|
|
|
|
Robuste par construction : si le nombre de sous-modèles disponibles ne
|
|
correspond pas à scaler.n_features_in_ (ex : DL entraîné mais
|
|
indisponible à l'inférence), on retombe sur une moyenne simple plutôt que
|
|
de lever une exception — CE FALLBACK EST IDENTIQUE PARTOUT, alors qu'avant
|
|
ahad_quant.py plantait dans ce cas alors que rl_env.py silencieusement
|
|
dégradait sans le signaler.
|
|
"""
|
|
raw_feats = np.asarray(raw_feats, dtype=np.float32)
|
|
if raw_feats.ndim == 1:
|
|
raw_feats = raw_feats.reshape(1, -1)
|
|
N = len(raw_feats)
|
|
|
|
if ensemble is None:
|
|
return np.column_stack([
|
|
np.full(N, 0.5, dtype=np.float32),
|
|
np.zeros(N, dtype=np.float32),
|
|
])
|
|
|
|
preds, names = [], []
|
|
|
|
lgbm = ensemble.get("lgbm")
|
|
xgb = ensemble.get("xgb")
|
|
rf = ensemble.get("rf")
|
|
|
|
if lgbm is not None:
|
|
preds.append(np.asarray(lgbm.predict(raw_feats), dtype=np.float32)); names.append("lgbm")
|
|
if xgb is not None:
|
|
preds.append(xgb.predict_proba(raw_feats)[:, 1].astype(np.float32)); names.append("xgb")
|
|
if rf is not None:
|
|
preds.append(rf.predict_proba(raw_feats)[:, 1].astype(np.float32)); names.append("rf")
|
|
|
|
if use_dl and has_dl_models(ensemble):
|
|
tft_p, tgru_p = _predict_dl_batch(ensemble, raw_feats)
|
|
if tft_p is not None:
|
|
preds.append(tft_p); names.append("tft")
|
|
if tgru_p is not None:
|
|
preds.append(tgru_p); names.append("tgru")
|
|
|
|
if not preds:
|
|
proba = np.full(N, 0.5, dtype=np.float32)
|
|
else:
|
|
meta = ensemble.get("meta")
|
|
scaler = ensemble.get("scaler")
|
|
stacked = np.column_stack(preds)
|
|
|
|
if (meta is not None and scaler is not None
|
|
and stacked.shape[1] == getattr(scaler, "n_features_in_", -1)):
|
|
stacked_s = scaler.transform(stacked)
|
|
proba = meta.predict_proba(stacked_s)[:, 1].astype(np.float32)
|
|
else:
|
|
if meta is not None and scaler is not None:
|
|
log.debug(
|
|
f"[ENS] méta-modèle attend {scaler.n_features_in_} entrées, "
|
|
f"{stacked.shape[1]} fournies ({names}) — fallback moyenne simple"
|
|
)
|
|
proba = stacked.mean(axis=1).astype(np.float32)
|
|
|
|
confidence = (np.abs(proba - 0.5) * 2.0).astype(np.float32)
|
|
return np.column_stack([proba, confidence])
|
|
|
|
|
|
def predict_ensemble_single(ensemble: dict, raw_feat_row: np.ndarray, history: np.ndarray = None) -> tuple:
|
|
"""
|
|
Prédit pour UNE SEULE observation (inférence live, bougie par bougie).
|
|
|
|
history : fenêtre récente (H, F) se terminant par raw_feat_row, utilisée
|
|
UNIQUEMENT si l'ensemble a des modèles séquentiels (has_dl).
|
|
Si absente ou trop courte, les composantes DL sont simplement
|
|
omises du stacking (comme tout sous-modèle absent) — jamais
|
|
d'exception.
|
|
|
|
Réutilise predict_ensemble_batch() en interne → live (1 ligne) et
|
|
entraînement/backtest (batch) ne PEUVENT PAS diverger silencieusement,
|
|
car c'est littéralement le même code qui tourne.
|
|
|
|
Retourne (proba, confidence) — deux floats Python.
|
|
"""
|
|
if ensemble is None:
|
|
return 0.5, 0.0
|
|
|
|
if history is not None and has_dl_models(ensemble):
|
|
window = np.asarray(history, dtype=np.float32)
|
|
if window.ndim == 1:
|
|
window = window.reshape(1, -1)
|
|
result = predict_ensemble_batch(ensemble, window, use_dl=True)
|
|
proba, confidence = float(result[-1, 0]), float(result[-1, 1])
|
|
return proba, confidence
|
|
|
|
feat_2d = np.asarray(raw_feat_row, dtype=np.float32).reshape(1, -1)
|
|
result = predict_ensemble_batch(ensemble, feat_2d, use_dl=False)
|
|
return float(result[0, 0]), float(result[0, 1])
|