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