""" AHAD QUANT — PPO Training Script (V6 — UNIFIED + RESUME) Entraîne l'agent RL avec Stable-Baselines3 PPO. NOUVEAUTÉS V6 + FIXES : [FIX BUG-RL-01] sys.stdout.reconfigure() wrappé dans try/except [FIX BUG-RL-02] Bloc argparse déplacé dans if __name__ == "__main__" [NEW] rl_progress.json — suivi de progression toutes les 10k steps [NEW] --resume flag — reprise automatique depuis dernier checkpoint [NEW] Drive backup automatique toutes les 100k steps [NEW] CHECKPOINT_FREQ réduit à 50k steps (plus de granularité) [NEW] ProgressCallback avec sauvegarde Drive intégrée ORDRE D'ENTRAÎNEMENT OBLIGATOIRE : 1. python download_data.py 2. python train.py → model_ensemble.pkl 3. python rl_train.py → rl_agent.zip Reprise après interruption : python rl_train.py --resume → reprend depuis dernier checkpoint Usage (local ou cloud) : !pip install stable-baselines3[extra] gymnasium shimmy !python rl_train.py !python rl_train.py --resume # si interruption Résultats : rl_agent.zip → modèle PPO final unifié rl_scaler.pkl → scaler z-score des features rl_checkpoints/ → checkpoints toutes les 50k steps + best_model.zip rl_train_log/ → logs TensorBoard rl_progress.json → progression (permet la reprise) """ import os, sys, pickle, json, time, shutil, glob import numpy as np # ─── FIX BUG-RL-01 : sys.stdout.reconfigure() incompatible avec certains terminaux ────────────── try: sys.stdout.reconfigure(encoding="utf-8", errors="replace") except AttributeError: pass # certains terminaux ne supportent pas reconfigure() # ─── Vérification dépendances ───────────────────────────────────────────────── try: from stable_baselines3 import PPO from stable_baselines3.common.env_util import make_vec_env from stable_baselines3.common.callbacks import ( EvalCallback, CheckpointCallback, BaseCallback ) from stable_baselines3.common.monitor import Monitor from stable_baselines3.common.vec_env import DummyVecEnv except ImportError: print("❌ stable-baselines3 non installé.") print(" pip install stable-baselines3[extra] gymnasium shimmy") sys.exit(1) import config from rl_env import AhadQuantForexEnv # ─── Vérification prérequis : ensemble doit être entraîné d'abord ───────────── ENS_PATH = getattr(config, "ENSEMBLE_MODEL_PATH", "model_ensemble.pkl") if not os.path.exists(ENS_PATH): print("=" * 65) print(" ❌ ERREUR : model_ensemble.pkl introuvable !") print(f" Chemin attendu : {ENS_PATH}") print(" Entraîner d'abord l'ensemble : python train.py") print("=" * 65) sys.exit(1) print(f" ✅ Ensemble trouvé : {ENS_PATH}") # ─── Fichier de progression (reprise après interruption) ───────────────────── PROGRESS_FILE = "rl_progress.json" def save_progress(steps_done, total_steps, last_checkpoint=None, start_time=None): """Sauvegarde la progression pour permettre la reprise automatique.""" elapsed = time.time() - start_time if start_time else 0 eta_sec = 0 if steps_done > 0 and elapsed > 0: rate = steps_done / elapsed eta_sec = (total_steps - steps_done) / rate if rate > 0 else 0 data = { "steps_done": steps_done, "total_steps": total_steps, "pct": round(steps_done / total_steps * 100, 2), "last_checkpoint": last_checkpoint, "timestamp": time.strftime("%Y-%m-%d %H:%M:%S"), "elapsed_min": round(elapsed / 60, 1), "eta_min": round(eta_sec / 60, 1), "status": "running", } with open(PROGRESS_FILE, "w") as f: json.dump(data, f, indent=2) def load_progress(): """Charge le fichier de progression s'il existe.""" if os.path.exists(PROGRESS_FILE): try: with open(PROGRESS_FILE) as f: return json.load(f) except Exception: return None return None def mark_completed(total_steps): """Marque l'entraînement comme terminé dans rl_progress.json.""" data = load_progress() or {} data["steps_done"] = total_steps data["total_steps"] = total_steps data["pct"] = 100.0 data["status"] = "completed" data["timestamp"] = time.strftime("%Y-%m-%d %H:%M:%S") with open(PROGRESS_FILE, "w") as f: json.dump(data, f, indent=2) # ─── Hyperparamètres PPO ────────────────────────────────────────────────────── PPO_CONFIG = dict( policy = "MlpPolicy", learning_rate = 3e-4, n_steps = 2048, batch_size = 256, n_epochs = 10, gamma = 0.99, gae_lambda = 0.95, clip_range = 0.2, ent_coef = 0.01, vf_coef = 0.5, max_grad_norm = 0.5, verbose = 1, tensorboard_log = "rl_train_log", policy_kwargs = dict( net_arch = [dict(pi=[256, 256, 128], vf=[256, 256, 128])], ), ) # ─── Config entraînement ───────────────────────────────────────────────────── DEFAULT_TOTAL_TIMESTEPS = 1_000_000 N_ENVS = 4 EVAL_FREQ = 20_000 CHECKPOINT_FREQ = 50_000 # [AMÉLIO] Réduit de 100k → 50k steps SAVE_PATH = "rl_agent" LOG_DIR = "rl_train_log" CHECKPOINT_DIR = "rl_checkpoints" DRIVE_BACKUP_DIR = "/content/drive/MyDrive/ahad_quant_checkpoints" DRIVE_BACKUP_EVERY = 100_000 # Backup Drive toutes les 100k steps TRAIN_PAIRS = [ "EURUSD", "GBPUSD", "USDJPY", "USDCHF", "AUDUSD", "USDCAD", "EURGBP", "EURJPY", ] EVAL_PAIRS = ["GBPJPY", "NZDUSD"] # ─── Callback métriques episodes ───────────────────────────────────────────── class EpisodeMetricsCallback(BaseCallback): def __init__(self, verbose=0): super().__init__(verbose) def _on_step(self) -> bool: for info in self.locals.get("infos", []): if "episode_metrics" in info: m = info["episode_metrics"] self.logger.record("ep/win_rate", m.get("ep/win_rate", 0)) self.logger.record("ep/total_trades", m.get("ep/total_trades", 0)) self.logger.record("ep/sharpe", m.get("ep/sharpe", 0)) return True # ─── Callback de progression et backup Drive ───────────────────────────────── class ProgressCallback(BaseCallback): """ [NEW] Sauvegarde rl_progress.json toutes les SAVE_EVERY steps. Permet la reprise automatique après une interruption imprévue. Effectue un backup Drive toutes les DRIVE_BACKUP_EVERY steps. """ def __init__(self, total_steps=1_000_000, save_every=10_000, drive_backup_every=100_000, verbose=0): super().__init__(verbose) self.total_steps = total_steps self.save_every = save_every self.drive_backup_every = drive_backup_every self._last_save = 0 self._last_drive_backup = 0 self._start_time = None def on_training_start(self, locals_, globals_): self._start_time = time.time() def _on_step(self) -> bool: ts = self.num_timesteps # Sauvegarde progression JSON if ts - self._last_save >= self.save_every: # Trouver le dernier checkpoint généré ckpts = sorted(glob.glob(f"{CHECKPOINT_DIR}/ppo_ahad_quant_unified_*.zip")) last_ckpt = ckpts[-1] if ckpts else None save_progress(ts, self.total_steps, last_ckpt, self._start_time) self._last_save = ts # Backup Drive automatique if ts - self._last_drive_backup >= self.drive_backup_every: self._backup_to_drive(ts) self._last_drive_backup = ts return True def _backup_to_drive(self, steps): """Copie les fichiers critiques vers Google Drive.""" drive_root = "/content/drive/MyDrive" if not os.path.exists(drive_root): return # Drive pas monté — non bloquant try: os.makedirs(DRIVE_BACKUP_DIR, exist_ok=True) files_to_backup = [ (PROGRESS_FILE, "rl_progress.json"), ("rl_scaler.pkl", "rl_scaler.pkl"), (f"{CHECKPOINT_DIR}/best_model.zip", "best_model.zip"), ] # Dernier checkpoint numéroté ckpts = sorted(glob.glob(f"{CHECKPOINT_DIR}/ppo_ahad_quant_unified_*.zip")) if ckpts: files_to_backup.append((ckpts[-1], "latest_checkpoint.zip")) copied = 0 for src, dst_name in files_to_backup: if os.path.exists(src): shutil.copy2(src, os.path.join(DRIVE_BACKUP_DIR, dst_name)) copied += 1 pct = steps / self.total_steps * 100 print(f"\n 💾 [Drive backup @ {steps:,} steps — {pct:.0f}%] {copied} fichiers → {DRIVE_BACKUP_DIR}") except Exception as e: print(f"\n ⚠️ Drive backup échoué (non bloquant) : {e}") # ─── Factory d'environnements ───────────────────────────────────────────────── def _make_env(pairs: list, seed: int = 0, ensemble_path: str = None): def _init(): env = AhadQuantForexEnv( data_dir = config.DATA_DIR, pairs = pairs, sl_pct = config.STOP_LOSS_PCT, tp_pct = config.TAKE_PROFIT_PCT, leverage = float(config.LEVERAGE), risk_per_trade = config.RISK_PER_TRADE, ensemble_path = ensemble_path, seed = seed, ) return Monitor(env) return _init # ─── Sauvegarde du scaler ───────────────────────────────────────────────────── def save_scaler(env: AhadQuantForexEnv, path: str = "rl_scaler.pkl"): scaler = { "mean": env._feat_mean, "std": env._feat_std, "pairs": env._pairs, } with open(path, "wb") as f: pickle.dump(scaler, f) print(f"[SCALER] Sauvegardé → {path}") return scaler # ─── Fonction de reprise — trouver le bon checkpoint ───────────────────────── def find_resume_checkpoint(): """ [NEW] Identifie le meilleur checkpoint depuis lequel reprendre. Priorité : last_checkpoint dans progress > best_model.zip > dernier numéroté """ prog = load_progress() if prog: last_ckpt = prog.get("last_checkpoint") if last_ckpt and os.path.exists(last_ckpt): return last_ckpt.replace(".zip", ""), prog best = f"{CHECKPOINT_DIR}/best_model.zip" if os.path.exists(best): return best.replace(".zip", ""), prog numbered = sorted(glob.glob(f"{CHECKPOINT_DIR}/ppo_ahad_quant_unified_*.zip")) if numbered: return numbered[-1].replace(".zip", ""), prog # Chercher aussi dans le backup Drive drive_ckpt = os.path.join(DRIVE_BACKUP_DIR, "latest_checkpoint.zip") drive_best = os.path.join(DRIVE_BACKUP_DIR, "best_model.zip") for drive_src in [drive_ckpt, drive_best]: if os.path.exists(drive_src): local_dst = os.path.join(CHECKPOINT_DIR, os.path.basename(drive_src)) os.makedirs(CHECKPOINT_DIR, exist_ok=True) shutil.copy2(drive_src, local_dst) print(f" ✅ Checkpoint récupéré depuis Drive : {os.path.basename(drive_src)}") return local_dst.replace(".zip", ""), prog return None, prog # ─── Entraînement principal ─────────────────────────────────────────────────── def train(total_steps: int = DEFAULT_TOTAL_TIMESTEPS, resume: bool = False): """ Entraîne ou reprend l'agent PPO. Args: total_steps: Nombre total de steps à entraîner. resume: Si True, tente de reprendre depuis le dernier checkpoint. """ print("=" * 65) print(" AHAD QUANT — RL Training Unifié (PPO V6)") print(" Ensemble supervisé intégré dans l'observation") if resume: print(" MODE : REPRISE depuis checkpoint") else: print(" MODE : Nouvel entraînement") print("=" * 65) os.makedirs(LOG_DIR, exist_ok=True) os.makedirs(CHECKPOINT_DIR, exist_ok=True) # ── Filtrer paires disponibles ───────────────────────────────────────────── available = [ p for p in TRAIN_PAIRS if os.path.exists(os.path.join(config.DATA_DIR, f"{p}_1h.json")) ] if not available: print(f"\n❌ Aucune donnée dans {config.DATA_DIR}. Exécuter download_data.py d'abord.") sys.exit(1) print(f"\n Paires entraînement : {available}") print(f" Ensemble : {ENS_PATH}") # ── Résolution checkpoint de reprise ────────────────────────────────────── resume_from = None steps_already = 0 remaining_steps = total_steps if resume: resume_from, prog = find_resume_checkpoint() if prog: steps_already = prog.get("steps_done", 0) remaining_steps = max(0, total_steps - steps_already) if steps_already >= total_steps: print(f"\n✅ Entraînement déjà complété ({steps_already:,} steps). Rien à faire.") return if resume_from: print(f"\n 🔄 Reprise depuis : {resume_from}.zip") print(f" ✅ Steps déjà effectués : {steps_already:,} / {total_steps:,} ({steps_already/total_steps*100:.1f}%)") print(f" ⏳ Steps restants : {remaining_steps:,}") else: print(" ⚠️ Aucun checkpoint trouvé — entraînement depuis zéro") resume = False print(f"\n Steps à entraîner : {remaining_steps:,}") print(f" Checkpoints toutes les {CHECKPOINT_FREQ:,} steps") print(f" Backup Drive toutes les {DRIVE_BACKUP_EVERY:,} steps\n") # ── Envs parallèles ──────────────────────────────────────────────────────── env_fns = [_make_env(available, seed=i, ensemble_path=ENS_PATH) for i in range(N_ENVS)] vec_env = DummyVecEnv(env_fns) # ── Env d'évaluation ────────────────────────────────────────────────────── eval_pairs = [p for p in EVAL_PAIRS if os.path.exists(os.path.join(config.DATA_DIR, f"{p}_1h.json"))] if not eval_pairs: eval_pairs = available[:2] print(f" Paires évaluation : {eval_pairs}") eval_env = Monitor(AhadQuantForexEnv( data_dir = config.DATA_DIR, pairs = eval_pairs, ensemble_path = ENS_PATH, seed = 999, )) # ── Scaler (toujours recalculé pour cohérence) ──────────────────────────── tmp_env = AhadQuantForexEnv( data_dir = config.DATA_DIR, pairs = available, ensemble_path = ENS_PATH, seed = 0, ) tmp_env.reset() save_scaler(tmp_env, "rl_scaler.pkl") obs_dim = tmp_env.obs_dim del tmp_env print(f"\n OBS_DIM : {obs_dim} dims") # ── Créer ou charger le modèle PPO ──────────────────────────────────────── if resume and resume_from and os.path.exists(f"{resume_from}.zip"): try: model = PPO.load( resume_from, env = vec_env, verbose = PPO_CONFIG["verbose"], tensorboard_log= PPO_CONFIG["tensorboard_log"], ) print(f"\n ✅ Modèle PPO chargé depuis checkpoint") except Exception as e: print(f"\n ⚠️ Échec chargement checkpoint : {e}") print(" → Redémarrage depuis zéro") model = PPO(env=vec_env, **PPO_CONFIG) resume = False steps_already = 0 remaining_steps = total_steps else: model = PPO(env=vec_env, **PPO_CONFIG) print(f"\n 🆕 Nouveau modèle PPO créé") # ── Callbacks ───────────────────────────────────────────────────────────── eval_cb = EvalCallback( eval_env, best_model_save_path = CHECKPOINT_DIR, log_path = LOG_DIR, eval_freq = EVAL_FREQ // N_ENVS, n_eval_episodes = 10, deterministic = True, verbose = 1, ) checkpoint_cb = CheckpointCallback( save_freq = CHECKPOINT_FREQ // N_ENVS, save_path = CHECKPOINT_DIR, name_prefix = "ppo_ahad_quant_unified", verbose = 1, ) metrics_cb = EpisodeMetricsCallback() progress_cb = ProgressCallback( total_steps = total_steps, save_every = 10_000, drive_backup_every = DRIVE_BACKUP_EVERY, ) # ── Lancement ───────────────────────────────────────────────────────────── print(f"\n Démarrage entraînement...\n") start = time.time() try: model.learn( total_timesteps = remaining_steps, callback = [eval_cb, checkpoint_cb, metrics_cb, progress_cb], reset_num_timesteps = not resume, # Ne pas reset si reprise progress_bar = True, ) except KeyboardInterrupt: print("\n\n⚠️ Entraînement interrompu manuellement (Ctrl+C).") print(f" Progression sauvegardée dans {PROGRESS_FILE}") print(" Relancer avec : python rl_train.py --resume") vec_env.close() eval_env.close() return except Exception as e: print(f"\n\n❌ Erreur pendant l'entraînement : {e}") print(f" Progression sauvegardée dans {PROGRESS_FILE}") print(" Relancer avec : python rl_train.py --resume") vec_env.close() eval_env.close() raise # ── Sauvegarde finale ───────────────────────────────────────────────────── elapsed = time.time() - start print(f"\n ✅ Entraînement terminé en {elapsed/60:.1f} min") model.save(SAVE_PATH) print(f" 💾 Modèle sauvegardé → {SAVE_PATH}.zip") mark_completed(total_steps) print(f" 📝 Progression marquée 'completed' → {PROGRESS_FILE}") vec_env.close() eval_env.close() # ── Résumé ──────────────────────────────────────────────────────────────── print("\n" + "=" * 65) print(" FICHIERS PRODUITS") print("=" * 65) for fname, desc in [ (f"{SAVE_PATH}.zip", "PPO unifié (cerveau final)"), ("rl_scaler.pkl", "Normalisation features"), (f"{CHECKPOINT_DIR}/best_model.zip", "Meilleur checkpoint"), (PROGRESS_FILE, "Fichier de progression"), (ENS_PATH, "Ensemble ML (feature PPO)"), ]: if os.path.exists(fname): size = os.path.getsize(fname) / 1024 / 1024 print(f" ✅ {fname:<45} {size:.1f} MB") else: print(f" ⚠️ {fname} — non trouvé") print() # ── Export unifié ───────────────────────────────────────────────────────── print("[Export] Modèle unifié...") try: from export_unified import export_unified unified_path = getattr(config, "UNIFIED_MODEL_PATH", "ahad_quant_unified.zip") export_unified( ppo_path = SAVE_PATH, ensemble_path = ENS_PATH, scaler_path = "rl_scaler.pkl", output_path = unified_path, ) except Exception as e: print(f" ⚠️ Export unifié échoué (non bloquant) : {e}") # ─── Fine-tuning classique ──────────────────────────────────────────────────── def fine_tune(model_path: str = "rl_agent", extra_steps: int = 200_000): """Fine-tune le modèle existant (appelé par auto_retrain.py ou CLI).""" print(f"[RL FINE-TUNE] Chargement {model_path}.zip...") available = [ p for p in TRAIN_PAIRS if os.path.exists(os.path.join(config.DATA_DIR, f"{p}_1h.json")) ] env_fns = [_make_env(available, seed=i, ensemble_path=ENS_PATH) for i in range(N_ENVS)] vec_env = DummyVecEnv(env_fns) model = PPO.load(model_path, env=vec_env) model.learn( total_timesteps = extra_steps, reset_num_timesteps = False, progress_bar = True, ) model.save(model_path) print(f"[RL FINE-TUNE] ✅ Sauvegardé → {model_path}.zip (+{extra_steps:,} steps)") vec_env.close() # ─── Fine-tuning avec replay ───────────────────────────────────────────────── def fine_tune_with_replay( model_path: str = "rl_agent", extra_steps: int = 200_000, replay_file: str = None, ) -> None: """Fine-tune PPO + injection d'épisodes réels depuis un fichier JSON.""" print(f"[RL FINE-TUNE+REPLAY] Chargement {model_path}.zip...") available = [ p for p in TRAIN_PAIRS if os.path.exists(os.path.join(config.DATA_DIR, f"{p}_1h.json")) ] env_fns = [_make_env(available, seed=i, ensemble_path=ENS_PATH) for i in range(N_ENVS)] vec_env = DummyVecEnv(env_fns) model = PPO.load(model_path, env=vec_env) if replay_file and os.path.exists(replay_file): try: with open(replay_file, "r") as f: episodes = json.load(f) print(f"[RL FINE-TUNE+REPLAY] {len(episodes)} épisodes réels chargés") injected = 0 for ep in episodes: try: obs = np.array(ep["obs"], dtype=np.float32) action = int(ep["action"]) reward = float(np.clip(ep["reward"], -1.0, 1.0)) done = bool(ep.get("done", False)) if hasattr(model, "rollout_buffer") and model.rollout_buffer is not None: buf = model.rollout_buffer if not buf.full: buf.add( obs.reshape(1, -1), np.array([[action]]), np.array([reward]), np.array([done]), np.zeros((1,), dtype=np.float32), np.zeros((1,), dtype=np.float32), ) injected += 1 except Exception: continue print(f"[RL FINE-TUNE+REPLAY] ✅ {injected}/{len(episodes)} épisodes injectés") except Exception as rep_err: print(f"[RL FINE-TUNE+REPLAY] ⚠️ Erreur injection (non bloquant) : {rep_err}") elif replay_file: print(f"[RL FINE-TUNE+REPLAY] ⚠️ Fichier replay introuvable : {replay_file}") model.learn( total_timesteps = extra_steps, reset_num_timesteps = False, progress_bar = True, ) model.save(model_path) print(f"[RL FINE-TUNE+REPLAY] ✅ Sauvegardé → {model_path}.zip (+{extra_steps:,} steps)") vec_env.close() # ─── FIX BUG-RL-02 : argparse déplacé dans if __name__ == "__main__" ────────── if __name__ == "__main__": import argparse parser = argparse.ArgumentParser(description="AHAD QUANT PPO Training") parser.add_argument("--finetune", action="store_true", help="Fine-tuner le modèle existant") parser.add_argument("--resume", action="store_true", help="[NEW] Reprendre depuis le dernier checkpoint") parser.add_argument("--steps", type=int, default=DEFAULT_TOTAL_TIMESTEPS, help=f"Nombre de steps (défaut: {DEFAULT_TOTAL_TIMESTEPS:,})") parser.add_argument("--model", type=str, default="rl_agent", help="Chemin du modèle pour fine-tune (sans .zip)") parser.add_argument("--replay", type=str, default=None, help="Fichier JSON d'épisodes réels pour fine-tune+replay") args = parser.parse_args() if args.finetune and args.replay: fine_tune_with_replay(args.model, args.steps, args.replay) elif args.finetune: fine_tune(args.model, args.steps) elif args.resume: train(total_steps=args.steps, resume=True) else: train(total_steps=args.steps, resume=False)