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

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