Three issues caused the dashboard to look frozen during the long Phase 1 exploration phase, even when 2+ runs had completed: 1. AI Thinking feed only got ONE entry (the phase-start banner) and then went silent until Phase 2. Phase 1 is LHS sampling so there are no AI calls — but we can still narrate observations rule-based. Added _narrate_phase1_run() that emits a per-result message tuned to what actually happened: 'PF 1.8, DD 14% — strong region, worth refining' (success), '$-734 loss with 191 trades — skipping' (warning), 'only 3 trades — params too restrictive' (info), etc. Mediocre runs are throttled to 1-in-5 to avoid feed spam 2. Parameter Changes panel showed 'No iterations yet — AI-driven param edits appear per iteration' which is misleading during Phase 1 (where there are no iterations YET, by design). Empty-state message now phase-aware: shows 'Exploration phase — changes appear in Phase 2' during phase1, 'Waiting for first AI iteration' if autonomous, or 'Random-neighbor refinement — enable Autonomous AI for AI changes' if not autonomous 3. AI Summary still said 'Waiting for optimization to start...' even while the optimization was actively running. New refreshAIWaitingState() updates the message based on (a) whether a run is active and (b) whether the API key is set. When running with no key, shows 'AI insights are off — no Anthropic API key set' with a 'Configure API key in Settings →' link Frontend now also tracks state.phase and state.autonomous from phase_start events so all empty-state messages stay in sync.
1295 lines
58 KiB
Python
1295 lines
58 KiB
Python
"""
|
||
optimizer/pipeline.py
|
||
The Smart 3-Phase Optimization Pipeline.
|
||
|
||
Replaces optimizer_loop.py as the core orchestration engine.
|
||
Receives a SessionConfig → runs Phase 1, 2, 3 → emits SocketIO events.
|
||
|
||
Architecture:
|
||
Phase 1: Broad Discovery (LHS samples, 20–30 runs)
|
||
Phase 2: Deep Refinement (neighbor search around top 3)
|
||
Phase 3: Validation (OOS backtest + sensitivity)
|
||
Decision: RECOMMENDED / RISKY / NOT_RELIABLE + .set file download
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import hashlib
|
||
import os
|
||
import random
|
||
import time
|
||
from datetime import datetime
|
||
from pathlib import Path
|
||
from typing import Optional
|
||
|
||
import yaml
|
||
from loguru import logger
|
||
|
||
from ea.registry import EARegistry
|
||
from ea.schema import ParameterSchema
|
||
from mt5.ini_builder import IniBuilder
|
||
from mt5.runner import MT5Runner
|
||
from mt5.report_parser import ReportParser
|
||
from data.models import Run
|
||
from data.store import DataStore
|
||
from reports.writer import ReportWriter
|
||
|
||
from optimizer.session_config import SessionConfig
|
||
from optimizer.lhs_sampler import LatinHypercubeSampler
|
||
from optimizer.result_ranker import ResultRanker, RankedResult
|
||
from optimizer.budget import BudgetManager
|
||
|
||
import pandas as pd
|
||
|
||
from analysis.equity_curve import EquityCurveAnalyzer
|
||
from analysis.time_performance import TimePerformanceAnalyzer
|
||
from analysis.ai_reasoner import AIReasoner
|
||
from analysis.ai_reasoner_config import load_api_key
|
||
|
||
BASE_DIR = Path(__file__).parent.parent
|
||
RUNS_DIR = BASE_DIR / "runs"
|
||
DB_PATH = BASE_DIR / "optimizer.db"
|
||
|
||
|
||
class OptimizationPipeline:
|
||
"""
|
||
3-phase autonomous optimization pipeline. Run in a background thread.
|
||
|
||
Usage:
|
||
pipeline = OptimizationPipeline("config.yaml", socketio, reports_dir)
|
||
pipeline.configure(session_config)
|
||
thread = threading.Thread(target=pipeline.run, daemon=True)
|
||
thread.start()
|
||
"""
|
||
|
||
def __init__(self, config_path: str, socketio, reports_dir: Path):
|
||
self.config_path = config_path
|
||
self.socketio = socketio
|
||
self.reports_dir = reports_dir
|
||
|
||
with open(config_path) as f:
|
||
self.cfg = yaml.safe_load(f)
|
||
|
||
# Runtime state
|
||
self.session: Optional[SessionConfig] = None
|
||
self.running = False
|
||
self._stop_flag = False
|
||
self._phase = "idle"
|
||
self._run_count = 0
|
||
self._total_runs = 0
|
||
self.best_result: Optional[RankedResult] = None
|
||
self.phase1_results: list[RankedResult] = []
|
||
self.phase2_results: list[RankedResult] = []
|
||
self.final_result: Optional[RankedResult] = None
|
||
self.verdict: Optional[str] = None
|
||
self.best_set_path: Optional[Path] = None
|
||
self.run_start_ts: Optional[float] = None
|
||
|
||
# AI & analysis state
|
||
self._ai_reasoner: Optional[AIReasoner] = None
|
||
self._ai_insights: list[dict] = []
|
||
self._run_findings: dict[str, list] = {}
|
||
self._run_insights: dict[str, dict] = {}
|
||
self._baseline_metrics: Optional[dict] = None
|
||
|
||
# Running best (updated each run so status can serve it live)
|
||
self._live_best: Optional[RankedResult] = None
|
||
# All completed runs (for history restoration)
|
||
self._completed_runs: list[dict] = []
|
||
|
||
# Live activity logs — persisted in-memory so a dashboard refresh during
|
||
# a run can replay them via /api/thinking, /api/param_changes, /api/validation.
|
||
self._thinking_log: list[dict] = []
|
||
self._param_changes: list[dict] = []
|
||
self._validation_log: list[dict] = []
|
||
self._early_term: Optional[dict] = None
|
||
|
||
# ── Public API ────────────────────────────────────────────────────────────
|
||
|
||
def configure(self, session: SessionConfig) -> None:
|
||
self.session = session
|
||
|
||
def reload_config(self) -> dict:
|
||
"""
|
||
Re-read config.yaml and apply runtime-mutable changes to the live
|
||
pipeline. Called from the Settings save handler so users can tune
|
||
AI model / timeout / thresholds mid-run without restarting.
|
||
|
||
Returns a dict of what actually changed, for the API response.
|
||
"""
|
||
try:
|
||
with open(self.config_path, encoding="utf-8") as f:
|
||
new_cfg = yaml.safe_load(f) or {}
|
||
except Exception as e:
|
||
return {"ok": False, "error": str(e)}
|
||
|
||
changed = {}
|
||
# AI settings — model, timeout, enabled — applied to the live reasoner
|
||
old_ai = self.cfg.get("ai", {})
|
||
new_ai = new_cfg.get("ai", {})
|
||
if self._ai_reasoner is not None:
|
||
for fld in ("model", "timeout_seconds", "enabled"):
|
||
if old_ai.get(fld) != new_ai.get(fld):
|
||
changed[f"ai.{fld}"] = new_ai.get(fld)
|
||
# Apply to live reasoner
|
||
try:
|
||
if "model" in new_ai and new_ai["model"]:
|
||
self._ai_reasoner.MODEL = new_ai["model"]
|
||
if "timeout_seconds" in new_ai:
|
||
self._ai_reasoner.TIMEOUT = int(new_ai["timeout_seconds"])
|
||
# API key swap: if a new full key was saved, rebuild reasoner
|
||
old_key = (old_ai.get("anthropic_api_key") or "").strip()
|
||
new_key = (new_ai.get("anthropic_api_key") or "").strip()
|
||
if new_key and new_key != old_key and not new_key.startswith("${") and len(new_key) >= 30:
|
||
self._ai_reasoner = AIReasoner(api_key=new_key)
|
||
changed["ai.api_key"] = "rotated"
|
||
except Exception as e:
|
||
logger.warning(f"reload_config: AI hot-reload failed: {e}")
|
||
|
||
# Threshold changes — applied to ResultRanker on next make_result.
|
||
# The thresholds module reads from config at import; we surface the
|
||
# change so the UI can confirm.
|
||
old_thr = self.cfg.get("thresholds", {})
|
||
new_thr = new_cfg.get("thresholds", {})
|
||
for fld in ("min_trades", "min_profit_factor", "min_calmar"):
|
||
if old_thr.get(fld) != new_thr.get(fld):
|
||
changed[f"thresholds.{fld}"] = new_thr.get(fld)
|
||
|
||
self.cfg = new_cfg
|
||
if changed:
|
||
self._log("info", f"⚙ Settings hot-reloaded: {', '.join(changed.keys())}")
|
||
self._emit("settings_reloaded", {"changed": changed})
|
||
return {"ok": True, "changed": changed}
|
||
|
||
def stop(self) -> None:
|
||
self._stop_flag = True
|
||
self._emit("status_change", {"state": "stopping"})
|
||
self._emit_early_termination(
|
||
reason_code="user_stop",
|
||
message="User requested stop. Finishing current run and exiting.",
|
||
details={"phase": self._phase},
|
||
)
|
||
|
||
def get_status(self) -> dict:
|
||
elapsed = int(time.time() - self.run_start_ts) if self.run_start_ts else 0
|
||
# Use final_result if set, otherwise best seen so far during the run
|
||
best = self.final_result or self._live_best
|
||
return {
|
||
"state": "running" if self.running else "idle",
|
||
"phase": self._phase,
|
||
"run_count": self._run_count,
|
||
"total_runs": self._total_runs,
|
||
"best_score": round(best.score, 4) if best else 0.0,
|
||
"verdict": self.verdict,
|
||
"elapsed_s": elapsed,
|
||
"ea_name": self.session.ea_name if self.session else "",
|
||
"symbol": self.session.symbol if self.session else "",
|
||
"timeframe": self.session.timeframe if self.session else "",
|
||
"has_insight": bool(self._ai_insights),
|
||
"insight_count": len(self._ai_insights),
|
||
"latest_insight": self._ai_insights[-1] if self._ai_insights else None,
|
||
"best_net_profit": round(best.net_profit, 2) if best else None,
|
||
"best_calmar": round(best.calmar, 3) if best else None,
|
||
"best_pf": round(best.profit_factor, 3) if best else None,
|
||
"best_win_rate": round(best.win_rate * 100, 1) if best else None,
|
||
"best_max_drawdown": round(best.max_drawdown * 100, 2) if best else None,
|
||
"best_total_trades": best.total_trades if best else None,
|
||
}
|
||
|
||
# ── Main entry point ──────────────────────────────────────────────────────
|
||
|
||
def run(self) -> None:
|
||
self.running = True
|
||
self._stop_flag = False
|
||
self.run_start_ts = time.time()
|
||
|
||
try:
|
||
self._run_pipeline()
|
||
except Exception as e:
|
||
logger.exception(f"Pipeline crashed: {e}")
|
||
self._emit("error", {"msg": f"Pipeline error: {e}"})
|
||
finally:
|
||
self.running = False
|
||
self._phase = "idle"
|
||
self._emit("status_change", {"state": "idle"})
|
||
|
||
def _run_pipeline(self) -> None:
|
||
cfg = self.session
|
||
self._total_runs = cfg.total_budget_runs
|
||
|
||
self._emit("status_change", {"state": "running", "phase": "setup"})
|
||
self._log("info", f"🚀 Smart Optimizer started — {cfg.ea_name} | {cfg.symbol} | {cfg.timeframe}")
|
||
self._log("info", f"📋 Budget: {cfg.budget_minutes} min | Samples: {cfg.phase1_samples} Phase1 + {cfg.phase2_samples} Phase2 + {cfg.phase3_samples} Phase3")
|
||
|
||
# ── Build components ─────────────────────────────────────────────────
|
||
reg = EARegistry(self.config_path)
|
||
profile = reg.get(cfg.ea_name)
|
||
schema = reg.get_schema(profile)
|
||
|
||
# Apply user's param selection if specified
|
||
if cfg.selected_params:
|
||
for p in schema.parameters.values():
|
||
if p.type != "fixed":
|
||
p.optimize = p.name in cfg.selected_params
|
||
|
||
builder = IniBuilder(self.config_path, schema=schema)
|
||
runner = MT5Runner(self.config_path)
|
||
parser = ReportParser()
|
||
store = DataStore(DB_PATH, RUNS_DIR)
|
||
writer = ReportWriter(self.reports_dir)
|
||
sampler = LatinHypercubeSampler(seed=int(time.time()))
|
||
ranker = ResultRanker(weights=cfg.scoring_weights)
|
||
budget = BudgetManager(cfg.budget_minutes)
|
||
|
||
# Initialize AI reasoning layer (model + timeout from config)
|
||
api_key = load_api_key(self.config_path)
|
||
ai_cfg = (self.cfg.get("ai") or {})
|
||
model = ai_cfg.get("model") or "claude-opus-4-7"
|
||
self._ai_reasoner = AIReasoner(api_key=api_key, model=model)
|
||
if "timeout_seconds" in ai_cfg:
|
||
try: self._ai_reasoner.TIMEOUT = int(ai_cfg["timeout_seconds"])
|
||
except Exception: pass
|
||
# Stream Claude's reasoning tokens to the dashboard as they arrive.
|
||
# The Live AI Thinking Feed listens for `ai_thinking_chunk` events.
|
||
try:
|
||
self._ai_reasoner.set_stream_callback(self._stream_ai_chunk)
|
||
except Exception:
|
||
pass
|
||
self._ai_insights.clear()
|
||
self._run_findings.clear()
|
||
self._run_insights.clear()
|
||
# Reset live-activity logs for the new run
|
||
self._thinking_log = []
|
||
self._param_changes = []
|
||
self._validation_log = []
|
||
self._early_term = None
|
||
|
||
budget.start()
|
||
|
||
# ── Phase 1: Broad Discovery ─────────────────────────────────────────
|
||
if self._stop_flag:
|
||
return
|
||
|
||
self._phase = "phase1"
|
||
self._emit("phase_start", {"phase": "phase1", "total": cfg.phase1_samples})
|
||
self._emit_thinking(
|
||
f"Starting broad discovery with {cfg.phase1_samples} Latin-Hypercube samples. "
|
||
f"Scanning parameter space to find profitable regions before focused refinement.",
|
||
kind="reasoning",
|
||
)
|
||
self._log("info", f"━━ Phase 1: Broad Discovery ({cfg.phase1_samples} configurations) ━━")
|
||
|
||
samples = sampler.sample(schema, cfg.phase1_samples)
|
||
phase1_raw: list[RankedResult] = []
|
||
|
||
for i, params in enumerate(samples):
|
||
if self._stop_flag:
|
||
break
|
||
|
||
run_id = f"p1_{i+1:02d}_{datetime.utcnow().strftime('%H%M%S')}"
|
||
self._log("info", f"[{i+1}/{cfg.phase1_samples}] Testing configuration {i+1}...")
|
||
|
||
t0 = time.time()
|
||
result = self._execute_run(
|
||
run_id, params, cfg.train_start, cfg.train_end,
|
||
"phase1", builder, runner, parser, store, writer, ranker, profile
|
||
)
|
||
budget.record_run(time.time() - t0)
|
||
phase1_raw.append(result)
|
||
self._run_count += 1
|
||
|
||
# Track live best and completed runs for status/history APIs
|
||
if self._live_best is None or result.score > self._live_best.score:
|
||
self._live_best = result
|
||
run_dict = self._make_run_dict(run_id, result, "phase1")
|
||
self._completed_runs.append(run_dict)
|
||
|
||
# Emit progress after each run
|
||
self._emit("run_complete", {
|
||
**run_dict,
|
||
"run_number": i + 1,
|
||
"total": cfg.phase1_samples,
|
||
"progress_pct": round((i + 1) / cfg.phase1_samples * 100),
|
||
"budget_summary": budget.summary(),
|
||
})
|
||
|
||
# Narrate every Phase 1 outcome so the AI Thinking feed stays alive
|
||
# even though no API call is made — this is observation, not AI reasoning
|
||
self._narrate_phase1_run(i + 1, cfg.phase1_samples, result)
|
||
|
||
if budget.is_exhausted():
|
||
self._log("warning", "⏱ Time budget reached during Phase 1")
|
||
break
|
||
|
||
# Rank Phase 1 results
|
||
self.phase1_results = ranker.rank(phase1_raw)
|
||
top5 = ranker.top_n(self.phase1_results, 5)
|
||
|
||
n_passing = sum(1 for r in self.phase1_results if r.passing)
|
||
self._log("info" if n_passing > 0 else "warning",
|
||
f"Phase 1 complete: {n_passing}/{len(self.phase1_results)} profitable configurations found"
|
||
)
|
||
|
||
# Emit Phase 1 summary for checkpoint UI
|
||
self._emit("phase1_complete", {
|
||
"total_tested": len(self.phase1_results),
|
||
"n_passing": n_passing,
|
||
"top_results": [self._result_to_dict(r) for r in top5],
|
||
})
|
||
|
||
if not top5:
|
||
self._emit("no_profitable_config", {
|
||
"msg": (
|
||
"Phase 1 found no profitable configuration. "
|
||
"Suggestions: try a different date range, check EA settings, "
|
||
"or try a different timeframe."
|
||
)
|
||
})
|
||
self._emit_early_termination(
|
||
reason_code="no_profit",
|
||
message="Optimization stopped early: Phase 1 found no profitable configuration.",
|
||
details={
|
||
"phase": "phase1",
|
||
"total_tested": len(self.phase1_results),
|
||
"suggestions": [
|
||
"Try a different date range",
|
||
"Check EA settings for issues",
|
||
"Try a different timeframe",
|
||
],
|
||
},
|
||
)
|
||
self._emit_thinking(
|
||
"No profitable configuration found across the broad scan. "
|
||
"Strategy is unstable under current EA settings — stopping optimization.",
|
||
kind="warning",
|
||
)
|
||
self._log("error", "❌ No profitable configuration found in Phase 1. Stopping.")
|
||
return
|
||
|
||
if self._stop_flag:
|
||
return
|
||
|
||
# ── Phase 2: Refinement (AI-Guided or Random Neighbor) ───────────────
|
||
if not budget.can_fit(3):
|
||
self._log("warning", "⏱ Not enough budget for Phase 2 — using Phase 1 winner directly")
|
||
self.final_result = top5[0]
|
||
|
||
elif cfg.autonomous_mode and self._ai_reasoner and self._ai_reasoner.enabled:
|
||
# ── AI-Guided Autonomous Loop ─────────────────────────────────────
|
||
self._phase = "phase2"
|
||
self._emit("phase_start", {
|
||
"phase": "phase2",
|
||
"total": cfg.autonomous_max_iterations,
|
||
"mode": "autonomous",
|
||
})
|
||
self._emit_thinking(
|
||
f"Phase 1 found {len(top5)} promising candidates. Best Calmar is "
|
||
f"{top5[0].calmar:.2f}. Switching to autonomous AI loop — I'll read each "
|
||
f"result, decide what parameters to change, and iterate toward the targets.",
|
||
kind="reasoning",
|
||
)
|
||
self._log("info",
|
||
f"━━ Phase 2: Autonomous AI Loop "
|
||
f"(up to {cfg.autonomous_max_iterations} iterations) ━━"
|
||
)
|
||
|
||
from optimizer.ai_guided_loop import AIGuidedLoop
|
||
loop = AIGuidedLoop(
|
||
pipeline=self, schema=schema, cfg=cfg,
|
||
builder=builder, runner=runner, parser=parser,
|
||
store=store, writer=writer, ranker=ranker,
|
||
profile=profile, budget=budget,
|
||
)
|
||
targets = {
|
||
"min_profit_factor": cfg.target_profit_factor,
|
||
"max_drawdown_pct": cfg.target_max_drawdown_pct,
|
||
"min_calmar": cfg.target_min_calmar,
|
||
}
|
||
loop.run(
|
||
seed_results=self.phase1_results,
|
||
max_iterations=cfg.autonomous_max_iterations,
|
||
targets=targets,
|
||
)
|
||
|
||
self.phase2_results = loop.all_results
|
||
all_candidates = [r for r in (self.phase1_results + loop.all_results) if r.passing]
|
||
all_ranked = ranker.rank(all_candidates) if all_candidates else ranker.rank(self.phase1_results)
|
||
self.final_result = all_ranked[0] if all_ranked else top5[0]
|
||
|
||
self._emit("phase2_complete", {
|
||
"best_run_id": self.final_result.run_id,
|
||
"best_score": round(self.final_result.score, 4),
|
||
"best_profit": round(self.final_result.net_profit, 2),
|
||
"best_calmar": round(self.final_result.calmar, 3),
|
||
"mode": "autonomous",
|
||
"iterations": len(loop.all_results),
|
||
})
|
||
self._log("info",
|
||
f"AI Loop complete. Best: {self.final_result.run_id} "
|
||
f"(PF={self.final_result.profit_factor:.2f}, "
|
||
f"profit=${self.final_result.net_profit:.0f}, "
|
||
f"calmar={self.final_result.calmar:.2f})"
|
||
)
|
||
|
||
else:
|
||
# ── Original Random-Neighbor Phase 2 ─────────────────────────────
|
||
self._phase = "phase2"
|
||
self._emit("phase_start", {"phase": "phase2", "total": cfg.phase2_samples})
|
||
self._emit_thinking(
|
||
f"Phase 2 starting in classic mode: sampling {cfg.phase2_samples} neighbors "
|
||
f"around the top {min(3, len(top5))} Phase-1 winners (no AI loop).",
|
||
kind="info",
|
||
)
|
||
self._log("info", f"━━ Phase 2: Deep Refinement (refining top {min(3, len(top5))} configs) ━━")
|
||
|
||
top3 = top5[:3]
|
||
phase2_raw = []
|
||
neighbors_per = cfg.phase2_samples // max(1, len(top3))
|
||
|
||
for rank_i, base_result in enumerate(top3):
|
||
if self._stop_flag:
|
||
break
|
||
|
||
self._log("info", f" Refining config #{rank_i+1}: {base_result.run_id}")
|
||
neighbors = sampler.sample_neighbors(
|
||
base_result.params, schema,
|
||
n_neighbors=neighbors_per, step_pct=0.20
|
||
)
|
||
|
||
for j, params in enumerate(neighbors):
|
||
if self._stop_flag or budget.is_exhausted():
|
||
break
|
||
|
||
run_id = f"p2_{rank_i+1}_{j+1:02d}_{datetime.utcnow().strftime('%H%M%S')}"
|
||
t0 = time.time()
|
||
result = self._execute_run(
|
||
run_id, params, cfg.train_start, cfg.train_end,
|
||
"phase2", builder, runner, parser, store, writer, ranker, profile
|
||
)
|
||
budget.record_run(time.time() - t0)
|
||
phase2_raw.append(result)
|
||
self._run_count += 1
|
||
|
||
run_dict_p2 = self._make_run_dict(run_id, result, "phase2")
|
||
self._completed_runs.append(run_dict_p2)
|
||
if self._live_best is None or result.score > self._live_best.score:
|
||
self._live_best = result
|
||
self._emit("run_complete", {
|
||
**run_dict_p2,
|
||
"progress_pct": round(self._run_count / self._total_runs * 100),
|
||
})
|
||
|
||
all_results = list(self.phase1_results) + ranker.rank(phase2_raw)
|
||
all_ranked = ranker.rank(
|
||
[r for r in all_results if r.passing]
|
||
or list(self.phase1_results)
|
||
)
|
||
self.phase2_results = ranker.rank(phase2_raw)
|
||
self.final_result = all_ranked[0] if all_ranked else top5[0]
|
||
|
||
self._emit("phase2_complete", {
|
||
"best_run_id": self.final_result.run_id,
|
||
"best_score": round(self.final_result.score, 4),
|
||
"best_profit": round(self.final_result.net_profit, 2),
|
||
"best_calmar": round(self.final_result.calmar, 3),
|
||
})
|
||
self._log("info",
|
||
f"Phase 2 complete. Best config: {self.final_result.run_id} "
|
||
f"(profit=${self.final_result.net_profit:.0f}, calmar={self.final_result.calmar:.2f})"
|
||
)
|
||
|
||
if self._stop_flag:
|
||
return
|
||
|
||
# ── Phase 3: Validation ───────────────────────────────────────────────
|
||
if not budget.can_fit(2):
|
||
self._log("warning", "⏱ Not enough budget for Phase 3 validation — skipping OOS test")
|
||
self._emit_early_termination(
|
||
reason_code="budget_exhausted",
|
||
message="Skipping validation: not enough time budget remaining.",
|
||
details={"phase": "phase3"},
|
||
)
|
||
self.verdict = "RISKY"
|
||
oos_result = None
|
||
else:
|
||
self._phase = "phase3"
|
||
self._emit("phase_start", {"phase": "phase3", "total": cfg.phase3_samples})
|
||
|
||
# Count the validation runs we intend to run so progress makes sense
|
||
opts = schema.optimizable()
|
||
plan_sens = bool(opts) and budget.can_fit(2)
|
||
planned_runs = 1 + (2 if plan_sens else 0) # 1 OOS + 2 sens
|
||
self._emit("validation_start", {
|
||
"best_run_id": self.final_result.run_id,
|
||
"planned_runs": planned_runs,
|
||
"oos_start": cfg.val_start,
|
||
"oos_end": cfg.val_end,
|
||
"train_start": cfg.train_start,
|
||
"train_end": cfg.train_end,
|
||
"sensitivity": plan_sens,
|
||
})
|
||
self._emit_thinking(
|
||
f"Entering validation phase. Testing best config {self.final_result.run_id} "
|
||
f"on out-of-sample data ({cfg.val_start} → {cfg.val_end}) + sensitivity checks.",
|
||
kind="reasoning",
|
||
)
|
||
self._log("info", "━━ Phase 3: Validation (out-of-sample + sensitivity) ━━")
|
||
|
||
# ── OOS test ──
|
||
oos_id = f"oos_{datetime.utcnow().strftime('%Y%m%d_%H%M%S')}"
|
||
self._log("info", f" OOS test: {cfg.val_start} → {cfg.val_end}")
|
||
self._emit("validation_run_start", {
|
||
"run_id": oos_id,
|
||
"kind": "oos",
|
||
"label": "Out-of-Sample",
|
||
"description": f"Re-testing best config on unseen data: {cfg.val_start} → {cfg.val_end}",
|
||
"index": 1,
|
||
"total": planned_runs,
|
||
"period_start": cfg.val_start,
|
||
"period_end": cfg.val_end,
|
||
"params": self.final_result.params,
|
||
})
|
||
self._emit_thinking(
|
||
"Running out-of-sample test: if the strategy holds up on data it wasn't "
|
||
"optimized on, the edge is real — not curve-fit noise.",
|
||
kind="hypothesis",
|
||
)
|
||
t0 = time.time()
|
||
oos_result = self._execute_run(
|
||
oos_id, self.final_result.params,
|
||
cfg.val_start, cfg.val_end,
|
||
"phase3_oos", builder, runner, parser, store, writer, ranker, profile
|
||
)
|
||
budget.record_run(time.time() - t0)
|
||
self._run_count += 1
|
||
|
||
oos_dict = self._make_run_dict(oos_id, oos_result, "phase3_oos")
|
||
self._completed_runs.append(oos_dict)
|
||
self._emit("run_complete", {**oos_dict, "progress_pct": round(self._run_count / self._total_runs * 100)})
|
||
self._emit("validation_run_complete", {
|
||
"run_id": oos_id,
|
||
"kind": "oos",
|
||
"label": "Out-of-Sample",
|
||
"index": 1,
|
||
"total": planned_runs,
|
||
"net_profit": round(oos_result.net_profit, 2),
|
||
"profit_factor": round(oos_result.profit_factor, 3),
|
||
"calmar": round(oos_result.calmar, 3),
|
||
"max_drawdown": round(oos_result.max_drawdown * 100, 2),
|
||
"total_trades": oos_result.total_trades,
|
||
"passing": oos_result.passing,
|
||
})
|
||
# Thinking narration on OOS outcome
|
||
oos_ratio = (oos_result.calmar / max(self.final_result.calmar, 0.001)) if oos_result.calmar else 0
|
||
if oos_result.net_profit > 0 and oos_ratio >= 0.50:
|
||
self._emit_thinking(
|
||
f"OOS profitable: ${oos_result.net_profit:.0f} with Calmar {oos_result.calmar:.2f} "
|
||
f"(≈{oos_ratio*100:.0f}% of training performance). The edge generalizes.",
|
||
kind="success",
|
||
)
|
||
else:
|
||
self._emit_thinking(
|
||
f"OOS weak: profit ${oos_result.net_profit:.0f}, Calmar {oos_result.calmar:.2f} — "
|
||
f"strategy likely overfit to the training period.",
|
||
kind="warning",
|
||
)
|
||
|
||
# ── Sensitivity test (2 runs: nudge top param up and down) ──
|
||
sens_results = []
|
||
if opts and budget.can_fit(2):
|
||
top_param = opts[0] # first optimizable param
|
||
self._emit_thinking(
|
||
f"Sensitivity check: nudging `{top_param.name}` ±20% to see if performance "
|
||
f"survives small parameter drift.",
|
||
kind="reasoning",
|
||
)
|
||
for i, direction in enumerate([1, -1], start=2):
|
||
if budget.is_exhausted() or self._stop_flag:
|
||
break
|
||
nudged = dict(self.final_result.params)
|
||
current = float(nudged.get(top_param.name, top_param.default))
|
||
span = float(top_param.max) - float(top_param.min)
|
||
nudged[top_param.name] = top_param.clamp(current + direction * span * 0.20)
|
||
|
||
sens_id = f"sens_{direction}_{datetime.utcnow().strftime('%H%M%S')}"
|
||
label = f"Sensitivity ({top_param.name} {'+' if direction > 0 else '-'}20%)"
|
||
self._emit("validation_run_start", {
|
||
"run_id": sens_id,
|
||
"kind": "sensitivity",
|
||
"label": label,
|
||
"description": f"Nudging `{top_param.name}` to {nudged[top_param.name]} "
|
||
f"(from {current}) to test parameter stability.",
|
||
"index": i,
|
||
"total": planned_runs,
|
||
"period_start": cfg.train_start,
|
||
"period_end": cfg.train_end,
|
||
"params": nudged,
|
||
})
|
||
t0 = time.time()
|
||
sr = self._execute_run(
|
||
sens_id, nudged, cfg.train_start, cfg.train_end,
|
||
"phase3_sens", builder, runner, parser, store, writer, ranker, profile
|
||
)
|
||
budget.record_run(time.time() - t0)
|
||
sens_results.append(sr)
|
||
self._run_count += 1
|
||
|
||
sens_dict = self._make_run_dict(sens_id, sr, "phase3_sens")
|
||
self._completed_runs.append(sens_dict)
|
||
self._emit("run_complete", {**sens_dict, "progress_pct": round(self._run_count / max(self._total_runs, 1) * 100)})
|
||
self._emit("validation_run_complete", {
|
||
"run_id": sens_id,
|
||
"kind": "sensitivity",
|
||
"label": label,
|
||
"index": i,
|
||
"total": planned_runs,
|
||
"net_profit": round(sr.net_profit, 2),
|
||
"profit_factor": round(sr.profit_factor, 3),
|
||
"calmar": round(sr.calmar, 3),
|
||
"max_drawdown": round(sr.max_drawdown * 100, 2),
|
||
"total_trades": sr.total_trades,
|
||
"passing": bool(sr.passing),
|
||
})
|
||
|
||
# Determine verdict
|
||
self.verdict = self._determine_verdict(self.final_result, oos_result, sens_results)
|
||
|
||
# Validation done — narrate the conclusion
|
||
self._emit("validation_done", {
|
||
"verdict": self.verdict,
|
||
"oos_passing": bool(oos_result and oos_result.passing),
|
||
"sens_passing": sum(1 for s in sens_results if s.passing),
|
||
"sens_total": len(sens_results),
|
||
})
|
||
verdict_narration = {
|
||
"RECOMMENDED": "Validation passed on all fronts. Strategy is robust and ready for deployment.",
|
||
"RISKY": "Validation mixed. Strategy may work but has stability concerns — proceed with care.",
|
||
"NOT_RELIABLE": "Validation failed. Strategy is not reliable — likely overfit or unstable.",
|
||
}.get(self.verdict, "Validation complete.")
|
||
self._emit_thinking(
|
||
verdict_narration,
|
||
kind=("success" if self.verdict == "RECOMMENDED"
|
||
else "warning" if self.verdict == "RISKY" else "warning"),
|
||
)
|
||
|
||
# ── Generate .set output ─────────────────────────────────────────────
|
||
self.best_set_path = self._write_set_file(self.final_result, schema, cfg)
|
||
|
||
# ── Final emit ───────────────────────────────────────────────────────
|
||
completion_payload = {
|
||
"verdict": self.verdict,
|
||
"best_run_id": self.final_result.run_id,
|
||
"score": round(self.final_result.score, 4),
|
||
"net_profit": round(self.final_result.net_profit, 2),
|
||
"calmar": round(self.final_result.calmar, 3),
|
||
"profit_factor": round(self.final_result.profit_factor, 3),
|
||
"win_rate": round(self.final_result.win_rate * 100, 1),
|
||
"max_drawdown": round(self.final_result.max_drawdown * 100, 2),
|
||
"total_trades": self.final_result.total_trades,
|
||
"oos_profit": round(oos_result.net_profit, 2) if oos_result else None,
|
||
"oos_calmar": round(oos_result.calmar, 3) if oos_result else None,
|
||
"set_file_url": f"/download_set/{self.final_result.run_id}" if self.best_set_path else None,
|
||
"total_runs": self._run_count,
|
||
"elapsed_min": round((time.time() - self.run_start_ts) / 60, 1),
|
||
}
|
||
self._emit("optimization_complete", completion_payload)
|
||
|
||
verdict_icon = {"RECOMMENDED": "✅", "RISKY": "⚠️", "NOT_RELIABLE": "❌"}.get(self.verdict, "?")
|
||
self._log("success" if self.verdict == "RECOMMENDED" else "warning",
|
||
f"{verdict_icon} VERDICT: {self.verdict} | "
|
||
f"Profit: ${self.final_result.net_profit:.0f} | "
|
||
f"Calmar: {self.final_result.calmar:.2f}"
|
||
)
|
||
|
||
# ── Webhook notification (Discord / Slack / generic) ─────────────────
|
||
try:
|
||
self._send_completion_webhook(completion_payload)
|
||
except Exception as e:
|
||
logger.debug(f"Webhook notification failed: {e}")
|
||
|
||
# ── Single run executor ───────────────────────────────────────────────────
|
||
|
||
def _execute_run(
|
||
self, run_id, params, period_start, period_end, phase,
|
||
builder, runner, parser, store, writer, ranker, profile
|
||
) -> RankedResult:
|
||
"""Execute one MT5 backtest and return a RankedResult."""
|
||
# Demo mode short-circuit — generate synthetic metrics so judges can see the
|
||
# AI loop without an MT5 install. Toggle with APEX_DEMO_MODE=1.
|
||
if os.environ.get("APEX_DEMO_MODE", "").strip() in ("1", "true", "yes"):
|
||
return self._execute_demo_run(run_id, params, phase, ranker)
|
||
|
||
run_dir = RUNS_DIR / run_id
|
||
run_dir.mkdir(parents=True, exist_ok=True)
|
||
|
||
try:
|
||
ini_path = builder.build(
|
||
run_id=run_id, params=params,
|
||
period_start=period_start, period_end=period_end,
|
||
output_dir=run_dir, phase=phase,
|
||
ea_file=profile.ex5_file,
|
||
ea_symbol=profile.symbol,
|
||
ea_timeframe=profile.timeframe,
|
||
)
|
||
|
||
run = Run(
|
||
run_id=run_id,
|
||
ea_name=profile.name, symbol=profile.symbol,
|
||
timeframe=profile.timeframe,
|
||
period_start=period_start, period_end=period_end,
|
||
params=params, phase=phase,
|
||
tester_model=self.cfg["mt5"]["tester_model"],
|
||
ini_snapshot=ini_path.read_text(),
|
||
)
|
||
store.save_run(run)
|
||
|
||
result = runner.run(
|
||
run_id, ini_path, run_dir / "report",
|
||
log_csv_search_dir=Path(self.cfg["mt5"]["mql5_files_path"]),
|
||
profile=profile,
|
||
)
|
||
|
||
if not result.success:
|
||
return ranker.make_result(run_id, params, phase, None,
|
||
error=result.error_message)
|
||
|
||
metrics, trades = parser.parse(result.report_xml, result.report_html)
|
||
if metrics is None:
|
||
return ranker.make_result(run_id, params, phase, None,
|
||
error="parse_failed")
|
||
|
||
metrics.run_id = run_id
|
||
|
||
# Run analysis & AI reasoning
|
||
findings = self._analyze_run(run_id, metrics, trades or [], params)
|
||
self._reason_about_run(run_id, metrics, findings, params)
|
||
|
||
trades_df = pd.DataFrame()
|
||
if trades:
|
||
try:
|
||
trades_df = pd.DataFrame([t.model_dump() for t in trades])
|
||
except Exception:
|
||
try:
|
||
trades_df = pd.DataFrame([vars(t) for t in trades])
|
||
except Exception:
|
||
pass
|
||
|
||
# Build RankedResult first so we have the ranked_score for summary.json
|
||
ranked = ranker.make_result(run_id, params, phase, metrics)
|
||
|
||
ai_insight = self._run_insights.get(run_id)
|
||
writer.write(
|
||
run_id, metrics, trades_df, findings, params,
|
||
phase=phase,
|
||
ranked_score=ranked.score,
|
||
ai_insight=ai_insight,
|
||
)
|
||
|
||
return ranked
|
||
|
||
except Exception as e:
|
||
logger.warning(f"[{run_id}] Run error: {e}")
|
||
return ranker.make_result(run_id, params, phase, None, error=str(e))
|
||
|
||
# ── Demo mode (synthetic backtest) ────────────────────────────────────────
|
||
|
||
def _execute_demo_run(self, run_id: str, params: dict, phase: str, ranker) -> RankedResult:
|
||
"""
|
||
Generate a synthetic, deterministic-but-realistic RankedResult from the
|
||
params hash. Lets the AI loop run end-to-end without MT5 installed.
|
||
|
||
The metrics improve as parameters approach a hidden "sweet spot" so the
|
||
AI can hill-climb in a way judges can observe. We also add small jitter
|
||
to look organic.
|
||
"""
|
||
from data.models import RunMetrics # local import to avoid cycle
|
||
|
||
# Hash params → stable seed per config
|
||
key = ",".join(f"{k}={v}" for k, v in sorted(params.items()))
|
||
seed = int(hashlib.md5(key.encode()).hexdigest()[:12], 16)
|
||
rng = random.Random(seed)
|
||
|
||
# Hidden sweet-spot hash drift: a deterministic "true score" between 0–1
|
||
# based on a smooth function of parameter values. Small param changes
|
||
# produce small score changes — that's what lets the AI hill-climb.
|
||
true_score = 0.0
|
||
for k, v in sorted(params.items()):
|
||
try:
|
||
fv = float(v)
|
||
# Map value into [-1,1] using a stable hash, multiply by a gentle
|
||
# bias toward "moderate" values (sweet spot is mid-range).
|
||
slot = (int(hashlib.md5(k.encode()).hexdigest()[:8], 16) % 100) / 100.0
|
||
normalized = (fv % 100) / 100.0
|
||
true_score += 1.0 - abs(normalized - slot)
|
||
except Exception:
|
||
true_score += 0.5
|
||
if params:
|
||
true_score /= len(params)
|
||
true_score = max(0.05, min(0.98, true_score))
|
||
|
||
# Add jitter for realism (±10%)
|
||
true_score *= rng.uniform(0.90, 1.10)
|
||
true_score = max(0.05, min(0.99, true_score))
|
||
|
||
# ── Regime failure ────────────────────────────────────────────────
|
||
# Real backtests produce occasional duds — bad luck on a market regime
|
||
# the parameters can't handle, or curve-fit overfit blowing up. Inject
|
||
# those so the verdict isn't always RECOMMENDED. Probability is higher
|
||
# in Phase 1 (pure exploration) and lower in Phase 2 (AI is improving).
|
||
regime_fail_prob = {
|
||
"phase1": 0.35, # 1 in 3 LHS samples lands in a bad spot
|
||
"phase2_ai": 0.10, # AI is improving, occasional misstep
|
||
"phase2": 0.20, # random neighbor search
|
||
"phase3_oos": 0.25, # the market the params never saw — sometimes brutal
|
||
"phase3_sens": 0.15, # nudged param → small chance of falling off the cliff
|
||
}.get(phase, 0.15)
|
||
|
||
regime_failed = rng.random() < regime_fail_prob
|
||
if regime_failed:
|
||
# Crush true_score so PF dives below 1, DD blows out, profit goes negative
|
||
true_score *= rng.uniform(0.10, 0.35)
|
||
|
||
# Project onto realistic metric ranges
|
||
profit_factor = round(0.7 + true_score * 1.8, 3) # 0.7–2.5
|
||
calmar = round(true_score * 1.4, 3) # 0.0–1.4
|
||
max_drawdown = round(28 - true_score * 22, 2) # 28%→6%
|
||
win_rate = round(40 + true_score * 25, 1) # 40%→65%
|
||
total_trades = int(80 + rng.random() * 220) # 80–300
|
||
net_profit = round((profit_factor - 1) * 5000 * (1 + rng.uniform(-0.2, 0.2)), 2)
|
||
|
||
# Out-of-sample is ALWAYS at least somewhat worse (reality bias).
|
||
# Roughly 30% of OOS runs degrade significantly (>40%) — that's what
|
||
# produces RISKY / NOT_RELIABLE verdicts in production.
|
||
if phase.startswith("phase3_oos"):
|
||
severe = rng.random() < 0.30
|
||
shrink = rng.uniform(0.35, 0.55) if severe else rng.uniform(0.75, 0.92)
|
||
profit_factor *= shrink
|
||
calmar *= shrink
|
||
net_profit *= shrink
|
||
if severe:
|
||
# Severe OOS failures also widen drawdown
|
||
max_drawdown = min(45.0, max_drawdown * rng.uniform(1.3, 1.8))
|
||
|
||
avg_trade = net_profit / max(total_trades, 1)
|
||
winners = int(total_trades * (win_rate / 100.0))
|
||
losers = max(1, total_trades - winners)
|
||
# Derive avg win/loss consistent with profit_factor: PF = (winners*avg_win)/(losers*|avg_loss|)
|
||
avg_loss = -abs(avg_trade) * (1 + 1.5 / max(profit_factor, 0.5))
|
||
avg_win = (profit_factor * losers * abs(avg_loss)) / max(winners, 1)
|
||
|
||
metrics = RunMetrics(
|
||
run_id=run_id,
|
||
net_profit=net_profit,
|
||
profit_factor=profit_factor,
|
||
calmar_ratio=calmar,
|
||
max_drawdown_pct=max_drawdown / 100.0,
|
||
max_drawdown_abs=round(net_profit * (max_drawdown / 100.0) * 1.2 + 200, 2),
|
||
win_rate=win_rate / 100.0,
|
||
total_trades=total_trades,
|
||
recovery_factor=round(profit_factor * 1.5, 2),
|
||
sharpe_ratio=round(true_score * 1.8, 2),
|
||
avg_win=round(avg_win, 2),
|
||
avg_loss=round(avg_loss, 2),
|
||
largest_loss=round(avg_loss * 3.5, 2),
|
||
expected_payoff=round(avg_trade, 2),
|
||
)
|
||
|
||
# Simulate per-run latency so the dashboard feels alive — not instant.
|
||
# Each run takes ~1.5s so a 10-iteration AI loop is ~15s total.
|
||
delay = float(os.environ.get("APEX_DEMO_RUN_SECONDS", "1.5"))
|
||
time.sleep(max(0.05, delay))
|
||
|
||
# Run AI reasoning if enabled — same flow as live mode
|
||
try:
|
||
self._reason_about_run(run_id, metrics, [], params)
|
||
except Exception:
|
||
pass
|
||
|
||
return ranker.make_result(run_id, params, phase, metrics)
|
||
|
||
# ── Verdict logic ─────────────────────────────────────────────────────────
|
||
|
||
def _determine_verdict(
|
||
self,
|
||
best: RankedResult,
|
||
oos: Optional[RankedResult],
|
||
sens: list[RankedResult],
|
||
) -> str:
|
||
"""
|
||
RECOMMENDED: IS profitable + OOS profitable + not fragile
|
||
RISKY: IS profitable but OOS weak OR fragile
|
||
NOT_RELIABLE: IS marginal or OOS loss
|
||
"""
|
||
if best.calmar < 0.1 or best.net_profit <= 0:
|
||
return "NOT_RELIABLE"
|
||
|
||
oos_ok = False
|
||
if oos and oos.net_profit > 0:
|
||
# OOS degradation: acceptable if OOS calmar ≥ 50% of IS calmar
|
||
oos_ratio = oos.calmar / max(best.calmar, 0.001)
|
||
oos_ok = oos_ratio >= 0.50
|
||
elif oos is None:
|
||
oos_ok = True # No OOS test — can't penalize
|
||
|
||
# Sensitivity: fragile if any nudge drops Calmar by >50%
|
||
fragile = any(
|
||
s.passing and s.calmar < best.calmar * 0.50
|
||
for s in sens
|
||
) if sens else False
|
||
|
||
if oos_ok and not fragile and best.calmar >= 0.30:
|
||
return "RECOMMENDED"
|
||
if oos_ok or (not fragile and best.calmar >= 0.20):
|
||
return "RISKY"
|
||
return "NOT_RELIABLE"
|
||
|
||
# ── .set file output ──────────────────────────────────────────────────────
|
||
|
||
def _write_set_file(
|
||
self, result: RankedResult, schema: ParameterSchema, cfg: SessionConfig
|
||
) -> Optional[Path]:
|
||
"""Write a clean .set file for MT5 import."""
|
||
try:
|
||
out_dir = self.reports_dir / result.run_id
|
||
out_dir.mkdir(parents=True, exist_ok=True)
|
||
out_path = out_dir / f"{cfg.ea_name}_optimized_{cfg.symbol}_{cfg.timeframe}.set"
|
||
|
||
header = (
|
||
f"Optimized by MT5 Smart Optimizer\n"
|
||
f"EA: {cfg.ea_name} | Symbol: {cfg.symbol} | TF: {cfg.timeframe}\n"
|
||
f"Training: {cfg.train_start} – {cfg.train_end}\n"
|
||
f"Verdict: {self.verdict}\n"
|
||
f"Net Profit: ${result.net_profit:.2f} | Calmar: {result.calmar:.2f} | "
|
||
f"Win Rate: {result.win_rate:.1f}%\n"
|
||
f"Generated: {datetime.utcnow().strftime('%Y-%m-%d %H:%M UTC')}"
|
||
)
|
||
|
||
content = schema.to_set_file(result.params, header_comment=header)
|
||
out_path.write_text(content, encoding="utf-8")
|
||
logger.info(f"Optimized .set file written: {out_path}")
|
||
return out_path
|
||
except Exception as e:
|
||
logger.warning(f"Could not write .set file: {e}")
|
||
return None
|
||
|
||
# ── Helpers ───────────────────────────────────────────────────────────────
|
||
|
||
def _make_run_dict(self, run_id: str, result: RankedResult, phase: str) -> dict:
|
||
"""Canonical run dict used for both _completed_runs and run_complete emits.
|
||
|
||
IMPORTANT — unit convention: win_rate and max_drawdown are emitted as
|
||
PERCENTAGES (0–100), not fractions, so the dashboard can render them
|
||
directly with `${val}%`. RankedResult stores these as fractions so we
|
||
scale here at the API boundary.
|
||
"""
|
||
d = {
|
||
"run_id": run_id,
|
||
"phase": phase,
|
||
"ts": datetime.utcnow().isoformat(),
|
||
"net_profit": round(result.net_profit, 2),
|
||
"calmar": round(result.calmar, 3),
|
||
"profit_factor": round(result.profit_factor, 3),
|
||
"win_rate": round(result.win_rate * 100, 1),
|
||
"max_drawdown": round(result.max_drawdown * 100, 2),
|
||
"total_trades": result.total_trades,
|
||
"passing": bool(result.passing),
|
||
"score": round(result.score, 4),
|
||
"params": result.params,
|
||
}
|
||
insight = self._run_insights.get(run_id)
|
||
if insight:
|
||
d["ai_insight"] = insight
|
||
return d
|
||
|
||
def _result_to_dict(self, r: RankedResult) -> dict:
|
||
d = {
|
||
"run_id": r.run_id,
|
||
"rank": r.rank,
|
||
"score": round(r.score, 4),
|
||
"net_profit": round(r.net_profit, 2),
|
||
"calmar": round(r.calmar, 3),
|
||
"profit_factor": round(r.profit_factor, 3),
|
||
"win_rate": round(r.win_rate * 100, 1), # fraction → %
|
||
"max_drawdown": round(r.max_drawdown * 100, 2), # fraction → %
|
||
"total_trades": r.total_trades,
|
||
"passing": bool(r.passing),
|
||
"params": r.params, # full param dict
|
||
"params_summary": self._params_summary(r.params),
|
||
}
|
||
# Attach AI insight if available
|
||
insight = self._run_insights.get(r.run_id)
|
||
if insight:
|
||
d["ai_insight"] = insight
|
||
return d
|
||
|
||
@staticmethod
|
||
def _params_summary(params: dict) -> str:
|
||
"""Show a few key params for display — works with any EA."""
|
||
parts = []
|
||
for k, v in list(params.items())[:8]:
|
||
short = k.replace("Inp", "").replace("inp", "")[:12]
|
||
parts.append(f"{short}={v}")
|
||
return " | ".join(parts[:4])
|
||
|
||
def _analyze_run(self, run_id: str, metrics, trades: list, params: dict) -> list:
|
||
"""Run Tier 1 analyzers on trade data. Always available from HTML report."""
|
||
findings = []
|
||
if not trades:
|
||
return findings
|
||
|
||
try:
|
||
trades_df = pd.DataFrame([t.model_dump() for t in trades])
|
||
except Exception:
|
||
try:
|
||
trades_df = pd.DataFrame([vars(t) for t in trades])
|
||
except Exception:
|
||
return findings
|
||
|
||
if trades_df.empty:
|
||
return findings
|
||
|
||
for analyzer_cls in (EquityCurveAnalyzer, TimePerformanceAnalyzer):
|
||
try:
|
||
az = analyzer_cls()
|
||
result = az.analyze(trades_df)
|
||
if isinstance(result, list):
|
||
findings.extend(result)
|
||
elif result is not None:
|
||
findings.append(result)
|
||
except Exception as e:
|
||
logger.debug(f"[{run_id}] {analyzer_cls.__name__} skipped: {e}")
|
||
|
||
self._run_findings[run_id] = findings
|
||
return findings
|
||
|
||
def _reason_about_run(
|
||
self, run_id: str, metrics, findings: list, params: dict
|
||
) -> Optional[dict]:
|
||
"""Call AI reasoner and emit insight via SocketIO."""
|
||
if not self._ai_reasoner or not self._ai_reasoner.enabled:
|
||
return None
|
||
try:
|
||
history = [
|
||
{
|
||
"run_id": r.run_id, "score": r.score,
|
||
"calmar": r.calmar, "pf": r.profit_factor, "phase": r.phase,
|
||
}
|
||
for r in (self.phase1_results + self.phase2_results)[-5:]
|
||
]
|
||
insight = self._ai_reasoner.analyze(
|
||
findings=findings,
|
||
metrics=metrics,
|
||
run_history=history,
|
||
current_params=params,
|
||
)
|
||
d = insight.to_dict()
|
||
d["run_id"] = run_id
|
||
self._ai_insights.append(d)
|
||
self._run_insights[run_id] = d
|
||
self._emit("ai_insight", d)
|
||
return d
|
||
except Exception as e:
|
||
logger.warning(f"[{run_id}] AI reasoning failed: {e}")
|
||
return None
|
||
|
||
def get_latest_insight(self) -> Optional[dict]:
|
||
"""Return the most recent AI insight."""
|
||
return self._ai_insights[-1] if self._ai_insights else None
|
||
|
||
def get_all_insights(self) -> list[dict]:
|
||
return list(self._ai_insights)
|
||
|
||
def _log(self, level: str, msg: str) -> None:
|
||
getattr(logger, level, logger.info)(msg)
|
||
self._emit("log", {"level": level, "msg": msg})
|
||
|
||
def _emit_thinking(
|
||
self,
|
||
msg: str,
|
||
kind: str = "info",
|
||
iteration: Optional[int] = None,
|
||
meta: Optional[dict] = None,
|
||
) -> None:
|
||
"""
|
||
Emit an AI-thinking stream event — distinct from system logs.
|
||
Shows up in the dashboard's "Live AI Thinking Feed".
|
||
|
||
kind: 'info' | 'reasoning' | 'decision' | 'warning' | 'success' | 'hypothesis'
|
||
"""
|
||
payload = {
|
||
"msg": msg,
|
||
"kind": kind,
|
||
"iteration": iteration,
|
||
"phase": self._phase,
|
||
"ts": datetime.utcnow().isoformat(),
|
||
}
|
||
if meta:
|
||
payload["meta"] = meta
|
||
# Persist for dashboard refresh — keep last 200 entries
|
||
self._thinking_log.append(payload)
|
||
if len(self._thinking_log) > 200:
|
||
self._thinking_log = self._thinking_log[-200:]
|
||
self._emit("ai_thinking", payload)
|
||
|
||
def _narrate_phase1_run(self, idx: int, total: int, result: RankedResult) -> None:
|
||
"""
|
||
Emit a short observation per Phase 1 run so the AI Thinking feed
|
||
stays alive during the LHS exploration phase. No API call — pure
|
||
rule-based narration of what just happened.
|
||
|
||
Cadence:
|
||
* Every passing result → success message
|
||
* Every clearly bad result → warning
|
||
* Once every 5 mediocre runs → progress beat
|
||
"""
|
||
try:
|
||
pf = float(getattr(result, "profit_factor", 0) or 0)
|
||
dd = float(getattr(result, "max_drawdown", 0) or 0) # fraction
|
||
tr = int(getattr(result, "total_trades", 0) or 0)
|
||
np_ = float(getattr(result, "net_profit", 0) or 0)
|
||
passing = bool(getattr(result, "passing", False))
|
||
|
||
tag = f"Run {idx}/{total}"
|
||
|
||
if tr == 0 or tr < 10:
|
||
msg = (
|
||
f"{tag}: only {tr} trades — params likely too restrictive "
|
||
f"(session window, signal threshold, or risk filter blocking entries)."
|
||
)
|
||
kind = "info"
|
||
elif passing and pf >= 1.5 and dd * 100 <= 15:
|
||
msg = (
|
||
f"{tag}: PF {pf:.2f}, DD {dd*100:.1f}%, {tr} trades — "
|
||
f"strong region, worth refining."
|
||
)
|
||
kind = "success"
|
||
elif passing:
|
||
msg = (
|
||
f"{tag}: PF {pf:.2f}, DD {dd*100:.1f}% — passes gates but "
|
||
f"not a standout. Logged."
|
||
)
|
||
kind = "info"
|
||
elif np_ < 0:
|
||
msg = (
|
||
f"{tag}: ${np_:.0f} loss with {tr} trades. Skipping this region."
|
||
)
|
||
kind = "warning"
|
||
elif dd * 100 > 25:
|
||
msg = (
|
||
f"{tag}: drawdown {dd*100:.1f}% too high — EA over-leveraged "
|
||
f"in this parameter region."
|
||
)
|
||
kind = "warning"
|
||
else:
|
||
# Mediocre — only emit one in five so we don't spam
|
||
if idx % 5 != 0:
|
||
return
|
||
msg = f"{tag}: PF {pf:.2f}, DD {dd*100:.1f}% — mediocre, moving on."
|
||
kind = "info"
|
||
|
||
self._emit_thinking(msg, kind=kind)
|
||
except Exception as e:
|
||
logger.debug(f"_narrate_phase1_run failed: {e}")
|
||
|
||
def _emit_early_termination(self, reason_code: str, message: str, details: dict = None) -> None:
|
||
"""
|
||
Surface an early stop to the user.
|
||
reason_code examples:
|
||
- 'no_profit' (Phase 1 found nothing)
|
||
- 'targets_met' (AI loop hit all quality targets)
|
||
- 'budget_exhausted'(time budget used up)
|
||
- 'user_stop' (user clicked Stop)
|
||
- 'stuck_escape' (optimizer stuck, bailing)
|
||
"""
|
||
payload = {
|
||
"reason": reason_code,
|
||
"message": message,
|
||
"phase": self._phase,
|
||
"ts": datetime.utcnow().isoformat(),
|
||
}
|
||
if details:
|
||
payload["details"] = details
|
||
self._early_term = payload # remember for refresh
|
||
self._emit("early_termination", payload)
|
||
|
||
def _stream_ai_chunk(self, evt: dict) -> None:
|
||
"""
|
||
Receive an event from AIReasoner streaming and forward to the
|
||
dashboard. Events: {event: 'start'|'delta'|'end'|'error', text?, error?}.
|
||
The frontend appends deltas to a single growing thinking-feed bubble.
|
||
"""
|
||
try:
|
||
self._emit("ai_thinking_chunk", {
|
||
"event": evt.get("event"),
|
||
"text": evt.get("text", ""),
|
||
"error": evt.get("error", ""),
|
||
"phase": self._phase,
|
||
"ts": datetime.utcnow().isoformat(),
|
||
})
|
||
except Exception:
|
||
pass
|
||
|
||
def _send_completion_webhook(self, payload: dict) -> None:
|
||
"""
|
||
POST a completion summary to a user-configured webhook URL. Supports
|
||
Discord and Slack styles natively; everything else gets the raw JSON.
|
||
Silent no-op if no URL is configured.
|
||
"""
|
||
notif = (self.cfg or {}).get("notifications", {}) or {}
|
||
url = (notif.get("webhook_url") or "").strip()
|
||
if not url:
|
||
return
|
||
|
||
verdict_emoji = {"RECOMMENDED": "✅", "RISKY": "⚠️", "NOT_RELIABLE": "❌"}.get(payload.get("verdict"), "📊")
|
||
title = f"{verdict_emoji} APEX optimization {payload.get('verdict','complete')}"
|
||
line = (
|
||
f"Best run **{payload.get('best_run_id','?')}** "
|
||
f"— PF {payload.get('profit_factor', 0):.2f} · "
|
||
f"Calmar {payload.get('calmar', 0):.2f} · "
|
||
f"DD {payload.get('max_drawdown', 0):.1f}% · "
|
||
f"Profit ${payload.get('net_profit', 0):,.0f} · "
|
||
f"{payload.get('total_trades', 0)} trades · "
|
||
f"{payload.get('elapsed_min', 0)} min"
|
||
)
|
||
|
||
# Detect style
|
||
style = (notif.get("webhook_style") or "auto").lower()
|
||
if style == "auto":
|
||
if "discord.com" in url or "discordapp.com" in url:
|
||
style = "discord"
|
||
elif "slack.com" in url or "hooks.slack" in url:
|
||
style = "slack"
|
||
else:
|
||
style = "generic"
|
||
|
||
try:
|
||
import requests
|
||
if style == "discord":
|
||
body = {"content": f"**{title}**\n{line}"}
|
||
elif style == "slack":
|
||
body = {"text": f"*{title}*\n{line}"}
|
||
else:
|
||
body = {"title": title, "summary": line, **payload}
|
||
requests.post(url, json=body, timeout=5)
|
||
self._log("info", f"📨 Webhook ({style}) notified.")
|
||
except Exception as e:
|
||
logger.warning(f"Webhook POST failed: {e}")
|
||
|
||
def _emit(self, event: str, data: dict = {}) -> None:
|
||
# Tee select live-activity events into in-memory logs so a dashboard
|
||
# refresh during the run can replay them via REST.
|
||
try:
|
||
if event == "param_changes":
|
||
self._param_changes.append(data)
|
||
if len(self._param_changes) > 50:
|
||
self._param_changes = self._param_changes[-50:]
|
||
elif event in ("validation_start", "validation_run_start", "validation_run_complete", "validation_done"):
|
||
self._validation_log.append({"event": event, **data})
|
||
if len(self._validation_log) > 40:
|
||
self._validation_log = self._validation_log[-40:]
|
||
except Exception:
|
||
pass
|
||
try:
|
||
self.socketio.emit(event, data)
|
||
except Exception as e:
|
||
logger.debug(f"Emit error ({event}): {e}")
|