Browse AI-generated trading strategies shared by the community. Fork, learn, and build on each other's work.
| Score▼ | Strategy | Author | Win Rate▼ | Return▼ | PF▼ | MDD▼ | Trades▼ | Actions | ||
|---|---|---|---|---|---|---|---|---|---|---|
|
🥇
|
Bollinger reversion
|
M
@malcolmtan
|
EURUSD | 1min | 47.9%56.2% | +1.53%+3.56% | 1.461.38 | 0.67%0.67% | 7164 |
|
# ╔══════════════════════════════════════════════════════════════╗
# ║ STRATEGY REQUEST LOG ║
# ╚══════════════════════════════════════════════════════════════╝
# Generated : 2026-05-25 02:29:29
# Model : XGBoost
# Feature Eng. : buy when price closes below the lower Bollinger Band(20,2) and RSI(14) < 35, exit at the middle band + Auto-add features: ON
# Signal / Entry : —
# Optimization : —
# Risk Mgmt : —
# Risk Filter : —
# ══════════════════════════════════════════════════════════════
# ============================================================
# Bollinger Band Mean-Reversion + RSI Filter (XGBoost, Sharpe)
# ============================================================
# SECTION 0 — IMPORTS & CONSTANTS
import numpy as np
import pandas as pd
# ── Inlined strategy_utils ──
"""
strategy_utils.py — Standard utility functions for generated strategies.
Claude imports these instead of writing boilerplate from scratch.
This ensures consistent behavior across all generated strategies.
"""
import numpy as np
import pandas as pd
from sklearn.preprocessing import LabelEncoder
# Max backtest window per timeframe. A finer timeframe over a longer window
# blows up the results dict / parquet load / Modal train time (the 2026-05-12
# OOM was a 1-min × multi-year sweep) — and a 1-min strategy gains nothing from
# 2 years of 1-min bars. Enforced HERE because every training path (UI / API /
# Modal) funnels through run_strategy → load_ohlc. Env-overridable so a future
# "max plan" / dedicated-server tier can lift it.
_TF_MAX_DAYS = {
"1min": 30,
"5min": 90,
"15min": 365,
"1h": 730,
}
def _fetch_ohlc_from_internal(symbol: str, tf: str, start: str, end: str):
"""Phase 3.2: fetch parquet bytes from Server A's /internal/ohlc endpoint
instead of reading a local file. Used inside Modal containers / Mac worker
pool (Phase 3.4) so every train sees the same source of truth as the chart.
Returns: pd.DataFrame (parquet decoded), or raises on any failure so the
caller can fall back / surface a clear error in the job.
"""
import hashlib as _hashlib, hmac as _hmac, io as _io, os as _os
import urllib.request as _ur, urllib.parse as _urp
base = (_os.environ.get("QM_INTERNAL_OHLC_BASE") or "").rstrip("/")
secret = (_os.environ.get("INTERNAL_WS_SECRET") or "").strip()
if not base:
raise RuntimeError("QM_INTERNAL_OHLC_BASE not set")
if not secret:
raise RuntimeError("INTERNAL_WS_SECRET not set")
msg = f"{symbol}|{tf}|{start}|{end}".encode("utf-8")
sig = _hmac.new(secret.encode("utf-8"), msg, _hashlib.sha256).hexdigest()
qs = _urp.urlencode({
"symbol": symbol, "tf": tf,
"start": start, "end": end, "sig": sig,
})
url = f"{base}/internal/ohlc?{qs}"
req = _ur.Request(url, headers={"User-Agent": "qm-worker/1.0"})
with _ur.urlopen(req, timeout=30) as resp:
if resp.status != 200:
raise RuntimeError(f"/internal/ohlc returned {resp.status}")
payload = resp.read()
print(f"[load_ohlc:internal] {symbol} {tf} fetched {len(payload)} bytes", flush=True)
return pd.read_parquet(_io.BytesIO(payload))
def _parse_symbol_tf_from_path(data_path: str):
"""Pull SYMBOL + TF out of a path like .../EURUSD_1min.parquet."""
import os as _os, re as _re
base = _os.path.basename(str(data_path))
m = _re.match(r"^([A-Z]{6})_(\d+min|\d+h)\.parquet$", base)
if not m:
return None, None
return m.group(1), m.group(2)
def load_ohlc(data_path, start_date="", end_date=""):
"""Load OHLC parquet, sort index, filter dates. Always returns consistent format.
The lower bound is clamped per timeframe (see _TF_MAX_DAYS) — a request for
more history than the cap silently starts later.
Phase 3.2: when env QM_USE_INTERNAL_OHLC=="1", fetch over HTTP from
Server A's /internal/ohlc endpoint instead of pd.read_parquet on a local
file (which on Modal is a stale Volume snapshot). The endpoint applies the
same day-cap, so the local cap-check below is a defensive no-op in that
path. Flag defaults to "0" → unchanged behavior.
Returns: (df, close, open_, high, low)
"""
import os as _os, re as _re
_use_internal = _os.environ.get("QM_USE_INTERNAL_OHLC", "0") == "1"
if _use_internal:
_sym, _tf = _parse_symbol_tf_from_path(data_path)
if not _sym or not _tf:
raise RuntimeError(
f"QM_USE_INTERNAL_OHLC=1 but DATA_PATH basename does not match "
f"SYMBOL_TF.parquet: {data_path}"
)
df = _fetch_ohlc_from_internal(_sym, _tf, start_date or "", end_date or "")
else:
df = pd.read_parquet(data_path)
df.index = pd.to_datetime(df.index)
df = df.sort_index()
# Per-timeframe window cap (timeframe inferred from the parquet filename).
_m = _re.search(r"_(\d+min|\d+h)\.parquet$", _os.path.basename(str(data_path)))
_tf = _m.group(1) if _m else None
_max_days = _TF_MAX_DAYS.get(_tf)
if _max_days and _max_days > 0 and len(df):
_env_override = _os.environ.get(f"QM_MAX_DAYS_{_tf.upper()}")
if _env_override and _env_override.isdigit():
_max_days = int(_env_override)
try:
_eff_end = pd.Timestamp(end_date) if end_date else df.index.max()
_eff_end = min(_eff_end, df.index.max())
_floor = _eff_end - pd.Timedelta(days=_max_days)
_req_start = pd.Timestamp(start_date) if start_date else df.index.min()
if _req_start < _floor:
print(f"[load_ohlc] {_tf} backtest window capped to {_max_days}d: "
f"start {_req_start.date()} -> {_floor.date()}", flush=True)
start_date = _floor
except Exception as _e:
print(f"[load_ohlc] window-cap check skipped ({_e})", flush=True)
if start_date:
df = df[df.index >= start_date]
if end_date:
df = df[df.index <= end_date]
return df, df["close"], df["open"], df["high"], df["low"]
def make_target(close, horizon=4):
"""Create target: direction N bars ahead. Default 4 bars = 1 hour on 15-min data.
Returns: target (pd.Series of -1, 0, 1)
"""
return np.sign(close.shift(-horizon) - close)
def split_data(df, target, feature_cols, train_split=0.7, validation_date=""):
"""Train/test split. Handles both ratio and date-based splits.
Drops NaN from target before splitting. Encodes labels to [0,1,2].
Returns: dict with keys:
X_train, X_test, y_train, y_test,
y_train_enc, y_test_enc, enc,
close_train, close_test,
split_idx, split_dt, n_train, n_test
"""
# Drop NaN from target
mask = target.notna()
df = df[mask].copy()
target = target[mask]
close = df["close"]
# Build feature matrix
X = df[feature_cols].copy()
X = X.bfill().ffill()
X = X.replace([np.inf, -np.inf], np.nan).fillna(0.0)
# Split
if validation_date:
split_idx = len(df[df.index <= validation_date])
else:
split_idx = int(len(df) * train_split)
split_idx = max(1, min(split_idx, len(df) - 1))
X_train = X.iloc[:split_idx]
X_test = X.iloc[split_idx:]
y_train = target.iloc[:split_idx]
y_test = target.iloc[split_idx:]
close_train = close.iloc[:split_idx]
close_test = close.iloc[split_idx:]
split_dt = str(df.index[split_idx])
# Label encoding — always fit on [-1, 0, 1]
enc = LabelEncoder()
enc.fit([-1, 0, 1])
y_train_enc = enc.transform(y_train)
y_test_enc = enc.transform(y_test)
return {
"df": df, "X_train": X_train, "X_test": X_test,
"y_train": y_train, "y_test": y_test,
"y_train_enc": y_train_enc, "y_test_enc": y_test_enc,
"enc": enc,
"close": close, "close_train": close_train, "close_test": close_test,
"split_idx": split_idx, "split_dt": split_dt,
"n_train": len(X_train), "n_test": len(X_test),
}
def compute_overlays(close, df_index):
"""Compute BB and MA overlays on full dataset. Always consistent.
Returns: (bb_dict, ma_dict)
"""
bb_mid = close.rolling(20).mean()
bb_std = close.rolling(20).std()
bb_upper = bb_mid + 2 * bb_std
bb_lower = bb_mid - 2 * bb_std
ma50 = close.rolling(50).mean()
ma100 = close.rolling(100).mean()
ma200 = close.rolling(200).mean()
def _safe(s):
s = s.reindex(df_index).bfill().ffill()
return [float(x) if (x is not None and not np.isnan(x) and not np.isinf(x)) else None
for x in s.values]
bb = {"upper": _safe(bb_upper), "mid": _safe(bb_mid), "lower": _safe(bb_lower)}
ma = {"ma50": _safe(ma50), "ma100": _safe(ma100), "ma200": _safe(ma200)}
return bb, ma
def run_backtest(signal, close, capital=10000, cost=2e-5):
"""Run backtest with transaction costs.
Uses price-based trade returns (same as webapp _compute_trades).
Signal 0 = hold (keep current position), not close.
Returns: dict with equity, trade_returns, long_returns, short_returns, bar_returns
"""
sig_arr = signal.values
price_arr = close.values
idx = signal.index
n = len(price_arr)
# Trade returns — price-based (matches webapp _compute_trades exactly)
trade_returns = []
long_returns = []
short_returns = []
trade_log = []
last_dir = None
entry_price = None
entry_bar = None
for i in range(n):
s = sig_arr[i]
c = price_arr[i]
if s != 0.0 and s != last_dir:
# Direction change — close previous trade, open new
if last_dir is not None and entry_price is not None and entry_price != 0:
ret = float(last_dir * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if last_dir == 1:
long_returns.append(ret)
else:
short_returns.append(ret)
trade_log.append({
"type": "Buy" if last_dir == 1 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[i]),
"entry_price": round(entry_price, 5),
"exit_price": round(c, 5),
"pnl": round(last_dir * (c - entry_price), 5),
"pnl_pct": round(ret * 100, 3),
"exit_reason": "signal",
})
entry_price = c
entry_bar = i
last_dir = s
# Close last open trade
if last_dir is not None and entry_price is not None and n > 0 and entry_price != 0:
c = price_arr[-1]
ret = float(last_dir * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if last_dir == 1:
long_returns.append(ret)
else:
short_returns.append(ret)
trade_log.append({
"type": "Buy" if last_dir == 1 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[-1]),
"entry_price": round(entry_price, 5),
"exit_price": round(c, 5),
"pnl": round(last_dir * (c - entry_price), 5),
"pnl_pct": round(ret * 100, 3),
"exit_reason": "end",
})
# Equity curve from trade returns
cumret = 1.0
equity_vals = np.full(n, float(capital))
trade_idx = 0
in_trade = False
t_entry_price = None
t_dir = None
for i in range(n):
s = sig_arr[i]
c = price_arr[i]
if s != 0.0 and s != t_dir:
if t_dir is not None and t_entry_price is not None and t_entry_price != 0:
t_ret = t_dir * (c - t_entry_price) / t_entry_price - cost
cumret *= (1 + t_ret)
t_entry_price = c
t_dir = s
equity_vals[i] = capital * cumret
# Bar returns for Sharpe
bar_returns = np.zeros(n)
for i in range(1, n):
if price_arr[i - 1] != 0 and last_dir is not None:
bar_returns[i] = sig_arr[i - 1] * (price_arr[i] - price_arr[i - 1]) / price_arr[i - 1] if sig_arr[i - 1] != 0 else 0.0
return {
"equity": pd.Series(equity_vals, index=close.index),
"trade_returns": trade_returns,
"long_returns": long_returns,
"short_returns": short_returns,
"bar_returns": bar_returns,
"trade_log": trade_log,
}
def compute_trade_stats(trades, capital=10000):
"""Single source of truth for trade statistics.
Every display path reads from this — no recomputation anywhere.
All values are rounded and JSON-safe (no inf/nan).
"""
if not trades:
return {"n": 0, "wins": 0, "losses": 0, "wr": 0, "avg": 0,
"best": 0, "worst": 0, "ret": 0, "np": 0, "mdd": 0,
"pf": 0, "rr": 0, "expect": 0}
w = [r for r in trades if r > 0]
l = [r for r in trades if r < 0]
cumret = 1.0
for r in trades:
cumret *= (1 + r)
net_p = capital * (cumret - 1)
# Max drawdown
eq = np.cumprod([1.0] + [1 + r for r in trades])
peak = np.maximum.accumulate(eq)
mdd = float(((eq - peak) / peak).min()) if len(eq) > 1 else 0.0
# Profit Factor
gross_w = sum(w) if w else 0
gross_l = abs(sum(l)) if l else 0
pf = gross_w / gross_l if gross_l > 0 else (9999.0 if gross_w > 0 else 0)
# Risk:Reward
avg_w = float(np.mean(w)) if w else 0
avg_l = abs(float(np.mean(l))) if l else 0
rr = avg_w / avg_l if avg_l > 0 else (9999.0 if avg_w > 0 else 0)
# Expectancy
expect = net_p / len(trades)
return {
"n": len(trades), "wins": len(w), "losses": len(l),
"wr": round(len(w) / len(trades), 4),
"avg": round(float(np.mean(trades)), 6),
"best": round(max(w), 6) if w else 0,
"worst": round(min(l), 6) if l else 0,
"ret": round(cumret - 1, 6),
"np": round(net_p, 2),
"mdd": round(mdd, 6),
"pf": round(pf, 2),
"rr": round(rr, 2),
"expect": round(expect, 2),
}
def compute_metrics(bt_result, close_test, capital=10000):
"""Compute all standard metrics from backtest result.
Uses trade-level compounding (same as webapp _trade_stats) for accuracy.
Returns: dict with total_ret, bh_ret, sharpe_strat, sharpe_bh, mdd, n_trades
"""
equity = bt_result["equity"]
trade_returns = bt_result["trade_returns"]
# Total return — trade-level compounding (matches webapp)
if trade_returns:
cumret = 1.0
for r in trade_returns:
cumret *= (1 + r)
total_ret = cumret - 1
else:
total_ret = 0.0
# Buy and hold
bh_equity = capital * (close_test / close_test.iloc[0])
bh_ret = (bh_equity.iloc[-1] - capital) / capital if capital != 0 else 0.0
# Sharpe ratio — trade-level (matches webapp: sqrt(252*26) annualization)
if len(trade_returns) >= 2 and float(np.std(trade_returns)) > 0:
sharpe_strat = float(np.mean(trade_returns) / np.std(trade_returns) * np.sqrt(252 * 26))
else:
sharpe_strat = 0.0
bh_rets = bh_equity.pct_change().dropna()
if len(bh_rets) > 1 and bh_rets.std() != 0:
sharpe_bh = float((bh_rets.mean() / bh_rets.std()) * np.sqrt(252 * 24 * 4))
else:
sharpe_bh = 0.0
# Max drawdown — trade-level (matches webapp)
if trade_returns:
eq = np.cumprod([1.0] + [1 + r for r in trade_returns])
peak = np.maximum.accumulate(eq)
mdd = float(((eq - peak) / peak).min()) if len(eq) > 1 else 0.0
else:
mdd = 0.0
return {
"total_ret": float(total_ret),
"bh_ret": float(bh_ret),
"sharpe_strat": float(sharpe_strat) if not np.isnan(sharpe_strat) else 0.0,
"sharpe_bh": float(sharpe_bh) if not np.isnan(sharpe_bh) else 0.0,
"mdd": float(mdd),
"n_trades": len(trade_returns),
}
# Diagnostics line/histogram series (equity / drawdown / rolling_acc / conf_hist)
# only feed the small Diagnostics charts — they're never used by the price chart
# or scroll-back. On a 1-min model trained over the (2.2-capped) window these are
# still ~30k points each; downsample to a visually-identical resolution before the
# dict leaves the trainer so it doesn't carry that into Server-A RAM / Postgres.
_RESULTS_SERIES_MAX = 5000
def _downsample_idx(n, cap=_RESULTS_SERIES_MAX):
"""Evenly-spaced index list spanning [0, n-1] (first+last always kept), or
None when no downsampling is needed (n <= cap)."""
if n <= cap:
return None
return np.unique(np.linspace(0, n - 1, cap).astype(int)).tolist()
def _take(arr, idx):
"""Subset a list by an index list (idx may be None → return arr unchanged)."""
if idx is None or not isinstance(arr, list):
return arr
return [arr[i] for i in idx]
# trade_log / train_trade_log are lists of per-trade dicts (display-only — the
# Trade Log tab). They scale with TRADE count, not bar count, so the bar-window
# cap (Phase 2.2) doesn't bound them — a degenerate near-every-bar model can put
# 10k+ trade dicts in the blob (>3 MB). Cap each (independently — a small-N model
# keeps every trade) to the most-recent N, recording `*_total` + `*_truncated`
# so the true count is still reported. Real strategies have far fewer than
# _TRADE_LOG_MAX trades, so this only ever bites pathological models.
_TRADE_LOG_MAX = 5000
def _cap_trade_log(tl):
"""Return (capped_list, original_len, was_truncated)."""
if not isinstance(tl, list) or len(tl) <= _TRADE_LOG_MAX:
return tl, (len(tl) if isinstance(tl, list) else 0), False
return tl[-_TRADE_LOG_MAX:], len(tl), True
def build_return_dict(split_result, bt_result, metrics, model, feature_cols,
signal_full, p_pos_test, p_neg_test, custom_figs=None,
bt_train_result=None, pre_stats=None):
"""Assemble the complete return dict. Handles ALL serialization.
Never returns Timestamps, numpy arrays, or non-JSON types.
Returns: JSON-safe dict with all required keys
"""
df = split_result["df"]
close = split_result["close"]
close_test = split_result["close_test"]
X_test = split_result["X_test"]
y_test = split_result["y_test"]
equity = bt_result["equity"]
bar_returns = bt_result["bar_returns"]
# OHLC
ohlc_dates = [str(x) for x in df.index.tolist()]
def _safe_list(arr):
return [float(x) if (x is not None and not np.isnan(x) and not np.isinf(x)) else None
for x in arr]
# Overlays
bb, ma = compute_overlays(close, df.index)
# Buy and hold equity
capital = equity.iloc[0] if len(equity) > 0 else 10000
bh_equity = capital * (close_test / close_test.iloc[0])
# Confusion matrix
from sklearn.metrics import confusion_matrix
pred_test = model.predict(X_test)
y_test_arr = np.asarray(y_test)
cm = confusion_matrix(y_test_arr, pred_test, labels=[-1, 0, 1])
# Rolling accuracy
sig_arr = signal_full.reindex(close_test.index).values
correct = pd.Series((pred_test == y_test_arr).astype(float), index=X_test.index)
active_test = pd.Series(sig_arr != 0, index=close_test.index) if len(sig_arr) == len(close_test) else pd.Series(True, index=close_test.index)
correct_active = correct.where(active_test, other=np.nan)
rolling_acc = correct_active.rolling(30, min_periods=1).mean()
# Feature importance
importances = model.feature_importances_
fi_pairs = sorted(zip(feature_cols, importances), key=lambda x: x[1])[-15:]
# Drawdown
rolling_max = equity.cummax()
drawdown = (equity - rolling_max) / rolling_max.replace(0, np.nan)
drawdown = drawdown.fillna(0.0)
# ── Downsample the Diagnostics-only series (see _downsample_idx) ──────────
_eq_dates = [str(x) for x in close_test.index.tolist()]
_eq_strat = _safe_list(equity.values)
_eq_bh = _safe_list(bh_equity.values)
_eq_idx = _downsample_idx(len(_eq_dates))
_eq_dates, _eq_strat, _eq_bh = _take(_eq_dates, _eq_idx), _take(_eq_strat, _eq_idx), _take(_eq_bh, _eq_idx)
_ra_dates = [str(x) for x in rolling_acc.index.tolist()]
_ra_vals = [float(x) if (not np.isnan(x) and not np.isinf(x)) else None for x in rolling_acc.values]
_ra_idx = _downsample_idx(len(_ra_dates))
_ra_dates, _ra_vals = _take(_ra_dates, _ra_idx), _take(_ra_vals, _ra_idx)
_dd_dates = [str(x) for x in drawdown.index.tolist()]
_dd_vals = _safe_list(drawdown.values)
_dd_idx = _downsample_idx(len(_dd_dates))
_dd_dates, _dd_vals = _take(_dd_dates, _dd_idx), _take(_dd_vals, _dd_idx)
_cp_pos = [float(x) for x in (p_pos_test.tolist() if hasattr(p_pos_test, 'tolist') else list(p_pos_test))]
_cp_neg = [float(x) for x in (p_neg_test.tolist() if hasattr(p_neg_test, 'tolist') else list(p_neg_test))]
_cp_pos = _take(_cp_pos, _downsample_idx(len(_cp_pos)))
_cp_neg = _take(_cp_neg, _downsample_idx(len(_cp_neg)))
# ── Trade logs — display-only (Trade Log tab); cap to most-recent N with a
# `_total` field so the true count is still reported (see _cap_trade_log).
# NB: ret_dist arrays are left FULL — a downstream path in callbacks.py
# recomputes n_trades/win-rate from len(ret_dist), so a sample would skew
# the displayed counts; they're small anyway and gzip handles them.
_tl_test, _tl_test_n, _tl_test_tr = _cap_trade_log(bt_result.get("trade_log", []))
_tl_tr, _tl_tr_n, _tl_tr_tr = _cap_trade_log(bt_train_result.get("trade_log", []) if bt_train_result else [])
return {
"ohlc": {
"dates": ohlc_dates,
"open": _safe_list(df["open"].values),
"high": _safe_list(df["high"].values),
"low": _safe_list(df["low"].values),
"close": _safe_list(df["close"].values),
},
"signals": {
"dates": [str(x) for x in signal_full.index.tolist()],
"values": [float(x) for x in signal_full.values],
},
"bb": bb,
"ma": ma,
"equity": {
"dates": _eq_dates,
"strategy": _eq_strat,
"bh": _eq_bh,
},
"feature_importance": {
"names": [p[0] for p in fi_pairs],
"values": [float(p[1]) for p in fi_pairs],
},
"conf_matrix": cm.tolist(),
"conf_hist": {
"p_pos": _cp_pos,
"p_neg": _cp_neg,
},
"rolling_acc": {
"dates": _ra_dates,
"values": _ra_vals,
},
"drawdown": {
"dates": _dd_dates,
"values": _dd_vals,
},
"ret_dist": [float(x) for x in bt_result["trade_returns"]],
"ret_dist_long": [float(x) for x in bt_result["long_returns"]],
"ret_dist_short": [float(x) for x in bt_result["short_returns"]],
"train_ret_dist": [float(x) for x in bt_train_result["trade_returns"]] if bt_train_result else [],
"train_ret_dist_long": [float(x) for x in bt_train_result["long_returns"]] if bt_train_result else [],
"train_ret_dist_short": [float(x) for x in bt_train_result["short_returns"]] if bt_train_result else [],
"trade_log": _tl_test,
"train_trade_log": _tl_tr,
"trade_log_total": _tl_test_n,
"train_trade_log_total": _tl_tr_n,
"trade_log_truncated": _tl_test_tr,
"train_trade_log_truncated": _tl_tr_tr,
**(pre_stats or {}),
"metrics": metrics,
"split_dt": split_result["split_dt"],
"split_idx": int(split_result["split_idx"]),
"n_train": int(split_result["n_train"]),
"n_test": int(split_result["n_test"]),
"feature_cols": list(feature_cols),
"custom_figs": custom_figs or [],
}
# ════════════════════════════════════════════════════════════════════════════
# STRATEGY FRAMEWORK v2 — Config-driven architecture
# Claude writes feature_engineering() + strategy_config(). Framework does rest.
# ════════════════════════════════════════════════════════════════════════════
import importlib
_MODEL_REGISTRY = {
"XGBClassifier": ("xgboost", "XGBClassifier"),
"RandomForestClassifier": ("sklearn.ensemble", "RandomForestClassifier"),
"GradientBoostingClassifier": ("sklearn.ensemble", "GradientBoostingClassifier"),
"LogisticRegression": ("sklearn.linear_model", "LogisticRegression"),
"ExtraTreesClassifier": ("sklearn.ensemble", "ExtraTreesClassifier"),
"AdaBoostClassifier": ("sklearn.ensemble", "AdaBoostClassifier"),
}
def _build_model_from_config(config, X_train, y_train_enc):
"""Build, fit, and wrap a model from strategy_config dict."""
model_type = config.get("model_type", "RandomForestClassifier")
model_params = dict(config.get("model_params", {}))
if model_type not in _MODEL_REGISTRY:
raise ValueError(f"Unknown model_type '{model_type}'. Valid: {list(_MODEL_REGISTRY.keys())}")
module_path, class_name = _MODEL_REGISTRY[model_type]
mod = importlib.import_module(module_path)
cls = getattr(mod, class_name)
# XGBoost defaults
if class_name == "XGBClassifier":
model_params.setdefault("use_label_encoder", False)
model_params.setdefault("eval_metric", "mlogloss")
model_params.setdefault("tree_method", "hist")
# Determinism > speed (2026-05-25). XGBoost hist with n_jobs=-1 is
# NON-reproducible even with random_state set — the parallel histogram
# gradient-sum order varies across threads, so the SAME code + data
# gives a slightly different model (and backtest) every run. Forcing
# single-thread makes training bit-reproducible so: (a) a user who
# copies a strategy and reruns it gets identical numbers, (b) the
# community "Live" score matches a redeploy, (c) "same code, different
# result" support reports go away. Cost: single-threaded XGB (a few
# seconds slower on large windows; hist is fast so it's minor). FORCED
# (not setdefault) so the guarantee can't be silently broken by a
# strategy passing n_jobs. Exact reproducibility holds within the
# platform (pinned versions / same Modal image); a user's own machine
# with different xgboost/numpy/CPU can still differ in low-order bits.
model_params["n_jobs"] = 1
# Common defaults
model_params.setdefault("random_state", 42)
from model_wrapper import ModelWrapper
clf = cls(**model_params)
clf.fit(X_train, y_train_enc)
enc = LabelEncoder()
enc.fit([-1, 0, 1])
return ModelWrapper(clf, original_classes=enc.classes_, n_features=X_train.shape[1])
def _generate_signals(model, X, threshold):
"""Framework-owned signal generation. Deterministic threshold logic."""
proba = model.predict_proba(X)
classes = list(model.classes_)
idx_pos = classes.index(1) if 1 in classes else None
idx_neg = classes.index(-1) if -1 in classes else None
p_pos = proba[:, idx_pos] if idx_pos is not None else np.zeros(len(X))
p_neg = proba[:, idx_neg] if idx_neg is not None else np.zeros(len(X))
signal_vals = np.zeros(len(X))
signal_vals = np.where(p_pos >= threshold, 1.0, signal_vals)
signal_vals = np.where(p_neg >= threshold, -1.0, signal_vals)
# Both exceed: pick stronger
both = (p_pos >= threshold) & (p_neg >= threshold)
signal_vals[both] = np.where(p_pos[both] >= p_neg[both], 1.0, -1.0)
return pd.Series(signal_vals, index=X.index), p_pos, p_neg
# ── Filter functions (all no-ops when config value is None) ──────────────
def _apply_direction_filter(signal, direction):
"""Zero out signals that don't match allowed direction."""
if direction is None or direction == "both":
return signal
s = signal.copy()
if direction == "long":
s[s < 0] = 0.0
elif direction == "short":
s[s > 0] = 0.0
return s
def _apply_session_filter(signal, index, session_hours):
"""Zero out signals outside session hours [start, end] UTC."""
if session_hours is None:
return signal
s = signal.copy()
start_h, end_h = session_hours[0], session_hours[1]
hours = index.hour
if start_h <= end_h:
mask = (hours >= start_h) & (hours < end_h)
else: # wrap around midnight, e.g. [22, 6]
mask = (hours >= start_h) | (hours < end_h)
s[~mask] = 0.0
return s
def _apply_atr_filter(signal, close, high, low, min_atr):
"""Zero out signals when NATR(14) is below threshold."""
if min_atr is None:
return signal
hl = high - low
hc = (high - close.shift(1)).abs()
lc = (low - close.shift(1)).abs()
tr = pd.concat([hl, hc, lc], axis=1).max(axis=1)
atr14 = tr.ewm(com=13, adjust=False).mean()
natr = atr14 / close.replace(0, np.nan)
s = signal.copy()
s[natr < min_atr] = 0.0
return s
def _apply_trend_filter(signal, close, trend_filter):
"""Only allow signals aligned with trend. e.g. 'sma_50': longs above SMA, shorts below."""
if trend_filter is None:
return signal
# Parse: "sma_50" → SMA with period 50
parts = trend_filter.lower().replace("-", "_").split("_")
if len(parts) >= 2 and parts[0] in ("sma", "ema"):
period = int(parts[1])
else:
return signal # unknown filter, skip
if parts[0] == "sma":
trend_line = close.rolling(period).mean()
else:
trend_line = close.ewm(span=period, adjust=False).mean()
s = signal.copy()
# Longs only above trend, shorts only below
s[(s > 0) & (close < trend_line)] = 0.0
s[(s < 0) & (close > trend_line)] = 0.0
return s
# ── run_backtest_v2: framework-owned SL/TP/cooldown/position management ──
def run_backtest_v2(signal, close, high, low, config, capital=10000, cost=2e-5):
"""Backtest with SL/TP/cooldown/direction handling built into the engine.
Unlike run_backtest (v1), this function handles position exits internally.
Returns: same dict shape as run_backtest()
"""
stop_loss = config.get("stop_loss")
take_profit = config.get("take_profit")
cooldown = config.get("cooldown", 0)
on_opposite = config.get("on_opposite", "reverse")
sig_arr = signal.values
close_arr = close.values
high_arr = high.values
low_arr = low.values
idx = signal.index
n = len(close_arr)
trade_returns = []
long_returns = []
short_returns = []
trade_log = []
equity_vals = np.full(n, float(capital))
cumret = 1.0
position = 0.0 # current direction: 1.0, -1.0, or 0.0 (flat)
entry_price = None
entry_bar = None # index into arrays for entry time
cooldown_remaining = 0
def _log_trade(exit_bar, exit_px, ret, reason):
trade_log.append({
"type": "Buy" if position == 1.0 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[exit_bar]),
"entry_price": round(entry_price, 5),
"exit_price": round(exit_px, 5),
"pnl": round(position * (exit_px - entry_price), 5),
"pnl_pct": round(ret * 100, 3),
"exit_reason": reason,
})
for i in range(n):
c = close_arr[i]
h = high_arr[i]
lo = low_arr[i]
s = sig_arr[i]
# 1. Check SL/TP if in trade
if position != 0.0 and entry_price is not None:
hit_sl = False
hit_tp = False
exit_price = None
if position == 1.0: # long
if stop_loss is not None and lo <= entry_price * (1 - stop_loss):
hit_sl = True
exit_price = entry_price * (1 - stop_loss)
elif take_profit is not None and h >= entry_price * (1 + take_profit):
hit_tp = True
exit_price = entry_price * (1 + take_profit)
else: # short
if stop_loss is not None and h >= entry_price * (1 + stop_loss):
hit_sl = True
exit_price = entry_price * (1 + stop_loss)
elif take_profit is not None and lo <= entry_price * (1 - take_profit):
hit_tp = True
exit_price = entry_price * (1 - take_profit)
if hit_sl or hit_tp:
ret = float(position * (exit_price - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, exit_price, ret, "SL" if hit_sl else "TP")
cumret *= (1 + ret)
position = 0.0
entry_price = None
entry_bar = None
cooldown_remaining = cooldown
equity_vals[i] = capital * cumret
continue
# 2. Cooldown
if cooldown_remaining > 0:
cooldown_remaining -= 1
equity_vals[i] = capital * cumret
continue
# 3. Signal processing
if s != 0.0:
if position == 0.0:
# Open new trade
position = s
entry_price = c
entry_bar = i
elif s != position:
# Opposite signal
if on_opposite == "reverse":
# Close current + open opposite
ret = float(position * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, c, ret, "signal")
cumret *= (1 + ret)
position = s
entry_price = c
entry_bar = i
else: # close_only
# Close current, go flat
ret = float(position * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, c, ret, "close_only")
cumret *= (1 + ret)
position = 0.0
entry_price = None
entry_bar = None
cooldown_remaining = cooldown
equity_vals[i] = capital * cumret
# Close last open trade at final close
if position != 0.0 and entry_price is not None and n > 0 and entry_price != 0:
c = close_arr[-1]
ret = float(position * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(n - 1, c, ret, "end")
cumret *= (1 + ret)
equity_vals[-1] = capital * cumret
# Bar returns for Sharpe (approximate)
bar_returns = np.zeros(n)
for i in range(1, n):
if close_arr[i - 1] != 0 and sig_arr[i - 1] != 0:
bar_returns[i] = sig_arr[i - 1] * (close_arr[i] - close_arr[i - 1]) / close_arr[i - 1]
return {
"equity": pd.Series(equity_vals, index=close.index),
"trade_returns": trade_returns,
"long_returns": long_returns,
"short_returns": short_returns,
"bar_returns": bar_returns,
"trade_log": trade_log,
}
# ── run_strategy: the v2 orchestrator ────────────────────────────────────
def run_strategy(feature_fn, config_fn, data_path, start_date="", end_date="",
validation_date="", train_split=0.7, register_model_fn=None):
"""Config-driven strategy execution. Claude writes feature_fn + config_fn,
framework does everything else.
Returns: results dict (same format as webapp expects)
"""
config = config_fn()
# Auto-correct SL/TP if Claude passed percentage instead of decimal
for _key in ("stop_loss", "take_profit"):
_val = config.get(_key)
if _val is not None and _val > 0.1: # >10% is almost certainly a percentage
config[_key] = _val / 100.0
print(f"[strategy] Auto-corrected {_key}: {_val} -> {config[_key]} (was percentage, converted to decimal)")
# 1. Load data
df, close, open_, high, low = load_ohlc(data_path, start_date, end_date)
# 2. Feature engineering (Claude's function)
df = feature_fn(df, close, open_, high, low)
close = df["close"]
open_ = df["open"]
high = df["high"]
low = df["low"]
# 3. Warm-up detection: drop rows where features have NaN BEFORE any fill
feature_cols = [c for c in df.columns if c not in ("open", "high", "low", "close")]
raw_nans = df[feature_cols].isna().any(axis=1)
valid_rows = ~raw_nans
if valid_rows.any():
first_valid = valid_rows.idxmax()
if raw_nans.loc[:first_valid].any():
df = df.loc[first_valid:].copy()
close = df["close"]
open_ = df["open"]
high = df["high"]
low = df["low"]
# 4. Target
horizon = config.get("target_horizon", 4)
target = make_target(close, horizon=horizon)
# 5. Split (ffill only within each partition — no bfill leak)
mask = target.notna()
df = df[mask].copy()
target = target[mask]
close = df["close"]
high = df["high"]
low = df["low"]
X = df[feature_cols].copy()
X = X.replace([np.inf, -np.inf], np.nan)
if validation_date:
split_idx = len(df[df.index <= validation_date])
else:
split_idx = int(len(df) * train_split)
split_idx = max(1, min(split_idx, len(df) - 1))
# ffill within train and test separately (no leak)
X_train = X.iloc[:split_idx].ffill().fillna(0.0)
X_test = X.iloc[split_idx:].ffill().fillna(0.0)
X = pd.concat([X_train, X_test])
y_train = target.iloc[:split_idx]
y_test = target.iloc[split_idx:]
close_train = close.iloc[:split_idx]
close_test = close.iloc[split_idx:]
high_test = high.iloc[split_idx:]
low_test = low.iloc[split_idx:]
enc = LabelEncoder()
enc.fit([-1, 0, 1])
y_train_enc = enc.transform(y_train)
y_test_enc = enc.transform(y_test)
split_dt = str(df.index[split_idx])
sp = {
"df": df, "X_train": X_train, "X_test": X_test,
"y_train": y_train, "y_test": y_test,
"y_train_enc": y_train_enc, "y_test_enc": y_test_enc,
"enc": enc,
"close": close, "close_train": close_train, "close_test": close_test,
"split_idx": split_idx, "split_dt": split_dt,
"n_train": len(X_train), "n_test": len(X_test),
}
# 6. Build model from config
model = _build_model_from_config(config, X_train, y_train_enc)
# 7. Generate signals
threshold = config.get("signal_threshold", 0.55)
signal_train, p_pos_train, p_neg_train = _generate_signals(model, X_train, threshold)
signal_test, p_pos_test, p_neg_test = _generate_signals(model, X_test, threshold)
# 8. Apply filters (order: direction → session → ATR → trend)
direction = config.get("direction", "both")
signal_test = _apply_direction_filter(signal_test, direction)
signal_train = _apply_direction_filter(signal_train, direction)
session_filter = config.get("session_filter")
signal_test = _apply_session_filter(signal_test, signal_test.index, session_filter)
signal_train = _apply_session_filter(signal_train, signal_train.index, session_filter)
min_atr = config.get("min_atr")
if min_atr is not None:
signal_test = _apply_atr_filter(signal_test, close_test, high_test, low_test, min_atr)
trend_filter = config.get("trend_filter")
if trend_filter is not None:
signal_test = _apply_trend_filter(signal_test, close_test, trend_filter)
signal_full = pd.concat([signal_train, signal_test])
# 9. Backtest with SL/TP/cooldown (test + train)
high_train = high.iloc[:split_idx]
low_train = low.iloc[:split_idx]
has_risk = (config.get("stop_loss") is not None or
config.get("take_profit") is not None or
config.get("cooldown", 0) > 0 or
config.get("on_opposite", "reverse") != "reverse")
if has_risk:
bt = run_backtest_v2(signal_test, close_test, high_test, low_test, config, capital=10000)
bt_train = run_backtest_v2(signal_train, close_train, high_train, low_train, config, capital=10000)
else:
bt = run_backtest(signal_test, close_test, capital=10000)
bt_train = run_backtest(signal_train, close_train, capital=10000)
# 10. Metrics
metrics = compute_metrics(bt, close_test, capital=10000)
# 11. Pre-compute all trade stats (single source of truth)
pre_stats = {
"train_stats": compute_trade_stats(bt_train.get("trade_returns", []), capital=10000),
"test_stats": compute_trade_stats(bt.get("trade_returns", []), capital=10000),
"long_stats": compute_trade_stats(bt.get("long_returns", []), capital=10000),
"short_stats": compute_trade_stats(bt.get("short_returns", []), capital=10000),
}
# 12. Register model
if register_model_fn is not None:
register_model_fn(model)
# 13. Build return dict
return build_return_dict(sp, bt, metrics, model, feature_cols,
signal_full, p_pos_test, p_neg_test, custom_figs=[],
bt_train_result=bt_train, pre_stats=pre_stats)
# ── End strategy_utils ──
DATA_PATH = '/root/Desktop/QuantifyMe/data/ohlc/AUDUSD_15min.parquet'
START_DATE = '2026-04-15'
END_DATE = '2026-05-25'
VALIDATION_DATE = ""
TRAIN_SPLIT = 0.7
# SECTION 1 — FEATURE ENGINEERING
def feature_engineering(df, close, open_, high, low):
# ── Bollinger Bands (20, 2) ──────────────────────────────────────────────
bb_period = 20
bb_std = 2.0
bb_mid = close.rolling(bb_period).mean()
bb_sigma = close.rolling(bb_period).std(ddof=0)
bb_upper = bb_mid + bb_std * bb_sigma
bb_lower = bb_mid - bb_std * bb_sigma
df["bb_mid"] = bb_mid
df["bb_upper"] = bb_upper
df["bb_lower"] = bb_lower
# %B — position of close within the band (0 = lower, 1 = upper)
bb_range = bb_upper - bb_lower
df["bb_pct_b"] = np.where(bb_range > 0, (close - bb_lower) / bb_range, 0.5)
# Bandwidth — normalised band width (regime filter)
df["bb_bandwidth"] = np.where(bb_mid > 0, bb_range / bb_mid, 0.0)
# Distance from each band (signed, normalised by sigma)
df["dist_lower"] = np.where(bb_sigma > 0, (close - bb_lower) / bb_sigma, 0.0)
df["dist_upper"] = np.where(bb_sigma > 0, (bb_upper - close) / bb_sigma, 0.0)
df["dist_mid"] = np.where(bb_sigma > 0, (close - bb_mid) / bb_sigma, 0.0)
# Below lower band flag
df["below_lower"] = np.where(close < bb_lower, 1, 0)
# Above upper band flag
df["above_upper"] = np.where(close > bb_upper, 1, 0)
# ── RSI (14) ─────────────────────────────────────────────────────────────
rsi_period = 14
delta = close.diff()
gain = delta.clip(lower=0)
loss = (-delta).clip(lower=0)
avg_gain = gain.ewm(com=rsi_period - 1, min_periods=rsi_period).mean()
avg_loss = loss.ewm(com=rsi_period - 1, min_periods=rsi_period).mean()
rs = np.where(avg_loss > 0, avg_gain / avg_loss, 100.0)
rsi = 100.0 - 100.0 / (1.0 + rs)
df["rsi"] = rsi
# RSI-derived flags and distances
df["rsi_oversold"] = np.where(rsi < 35, 1, 0)
df["rsi_overbought"] = np.where(rsi > 65, 1, 0)
df["rsi_dist_35"] = rsi - 35.0 # negative when oversold
df["rsi_dist_65"] = rsi - 65.0 # positive when overbought
df["rsi_norm"] = (rsi - 50.0) / 50.0 # centred, ±1 range
# ── Core entry condition features ────────────────────────────────────────
# Buy setup: close < lower BB AND RSI < 35
df["long_setup"] = np.where((close < bb_lower) & (rsi < 35), 1, 0)
# Sell setup: close > upper BB AND RSI > 65
df["short_setup"] = np.where((close > bb_upper) & (rsi > 65), 1, 0)
# ── ATR (14) — volatility context ────────────────────────────────────────
atr_period = 14
hl = high - low
hc = (high - close.shift(1)).abs()
lc = (low - close.shift(1)).abs()
tr = pd.concat([hl, hc, lc], axis=1).max(axis=1)
atr = tr.ewm(com=atr_period - 1, min_periods=atr_period).mean()
df["atr"] = atr
df["natr"] = np.where(close > 0, atr / close, 0.0)
# ── Momentum / Rate-of-Change ─────────────────────────────────────────────
for n in [1, 3, 5, 10]:
df[f"roc_{n}"] = np.where(
close.shift(n) > 0,
(close - close.shift(n)) / close.shift(n),
0.0
)
# ── EMA trend context (fast / slow) ──────────────────────────────────────
ema_fast = close.ewm(span=9, min_periods=9).mean()
ema_slow = close.ewm(span=21, min_periods=21).mean()
df["ema_fast"] = ema_fast
df["ema_slow"] = ema_slow
df["ema_diff"] = np.where(ema_slow > 0, (ema_fast - ema_slow) / ema_slow, 0.0)
df["ema_bull"] = np.where(ema_fast > ema_slow, 1, 0)
# SMA-50 trend filter helper (used by framework trend_filter)
df["sma_50"] = close.rolling(50).mean()
# ── Candle body & wick features ───────────────────────────────────────────
body = (close - open_).abs()
candle_rng = (high - low).replace(0, np.nan)
df["body_ratio"] = (body / candle_rng).fillna(0.0)
df["upper_wick"] = np.where(candle_rng.notna(), (high - close.clip(lower=open_)) / candle_rng.fillna(1), 0.0)
df["lower_wick"] = np.where(candle_rng.notna(), (close.clip(upper=open_) - low) / candle_rng.fillna(1), 0.0)
df["bull_candle"] = np.where(close > open_, 1, 0)
# ── Volume-like proxy — true range z-score ────────────────────────────────
tr_mean = tr.rolling(20).mean()
tr_std = tr.rolling(20).std(ddof=0).replace(0, np.nan)
df["tr_zscore"] = ((tr - tr_mean) / tr_std).fillna(0.0)
# ── Lagged RSI and %B (1, 2, 3 bars back) ────────────────────────────────
for lag in [1, 2, 3]:
df[f"rsi_lag{lag}"] = df["rsi"].shift(lag)
df[f"bb_pct_b_lag{lag}"] = df["bb_pct_b"].shift(lag)
# ── RSI slope ────────────────────────────────────────────────────────────
df["rsi_slope3"] = df["rsi"] - df["rsi"].shift(3)
# ── Mean-reversion proximity: how far price is from middle band ───────────
df["pct_to_mid"] = np.where(close > 0, (bb_mid - close) / close, 0.0)
# ── Fill any NaNs from warm-up ────────────────────────────────────────────
df = df.bfill().ffill()
return df
# SECTION 2 — STRATEGY CONFIG
def strategy_config():
return {
"title": "BB Mean-Reversion + RSI Oversold/Overbought (XGBoost)",
"model_type": "XGBClassifier",
"model_params": {
"n_estimators": 500,
"max_depth": 4,
"learning_rate": 0.03,
"subsample": 0.75,
"colsample_bytree": 0.70,
"min_child_weight": 5,
"gamma": 0.1,
"reg_alpha": 0.05,
"reg_lambda": 1.5,
"objective": "binary:logistic",
"random_state": 42,
"n_jobs": -1,
},
"signal_threshold": 0.55,
"direction": "both",
"stop_loss": 0.0010,
"take_profit": 0.0020,
"cooldown": 0,
"max_positions": 1,
"on_opposite": "reverse",
"session_filter": [7, 17],
"min_atr": 0.00005,
"trend_filter": None,
"target_horizon": 4,
"objective": (
"Maximize Sharpe ratio by exploiting Bollinger Band mean-reversion "
"with RSI confirmation. Entry conditions (close < lower BB, RSI < 35 "
"for longs; close > upper BB, RSI > 65 for shorts) are encoded as "
"features together with momentum, ATR volatility, candle structure, "
"and lagged indicators. XGBoost with strong regularisation "
"(reg_lambda=1.5, gamma=0.1, min_child_weight=5) and a low learning "
"rate avoids overfitting on the 6-week window. Session filter "
"[7,17] UTC targets liquid London/NY overlap, reducing noise. "
"TP:SL ratio of 2:1 supports positive expected value even at "
"moderate win rates, pushing Sharpe higher."
),
"notes": (
"Features: %B position, RSI (raw + flags + slope + lags), "
"EMA cross, ATR/NATR, ROC(1/3/5/10), candle body/wick ratios, "
"TR z-score, distance-to-midband, long/short setup flags. "
"Round-trip cost ~2e-5 is implicitly absorbed by the 10-pip TP target. "
"Cooldown=0 allows immediate re-entry after mean-reversion completes."
),
}
# ── Framework v2: auto-generated wrapper ──
def train_and_backtest():
_vd = VALIDATION_DATE if 'VALIDATION_DATE' in globals() else ''
_ts = TRAIN_SPLIT if 'TRAIN_SPLIT' in globals() else 0.7
return run_strategy(
feature_engineering, strategy_config,
DATA_PATH, START_DATE, END_DATE,
_vd, _ts,
register_model_fn=register_model
)
|
||||||||||
|
🥈
|
EUR/USD XGBoost RSI + Bollinger Bands Scalper
Maximize Sharpe ratio on EUR/USD 15-min data via XGBClassifier. Features: returns (1, 3, 5, 10, 20 bars), RSI 14, Bollinger Bands 20/2, SMAs…
NOISE · attempt 55
|
M
@malcolmtan
|
EURUSD | 1min | 61.1%60.9% | +6.61%-3.08% | 1.070.85 | 14.51%14.51% | 265192 |
|
QM_SWAP_SHORT_USD = 0.3
QM_SWAP_LONG_USD = -1.0
QM_COMMISSION_USD = 6.0
# ╔══════════════════════════════════════════════════════════════╗
# ║ STRATEGY REQUEST LOG ║
# ╚══════════════════════════════════════════════════════════════╝
# Generated : 2026-08-07 11:28:27
# Model : XGBoost
# Feature Eng. : Auto-add features: ON
# Signal / Entry : —
# Optimization : —
# Risk Mgmt : —
# Risk Filter : —
# ══════════════════════════════════════════════════════════════
# ============================================================
# QUANTIFY.ME TRADING STRATEGY — XGBoost RSI + Bollinger Bands
# ============================================================
import numpy as np
import pandas as pd
# ── Inlined strategy_utils ──
"""
strategy_utils.py — Standard utility functions for generated strategies.
Claude imports these instead of writing boilerplate from scratch.
This ensures consistent behavior across all generated strategies.
"""
import numpy as np
import pandas as pd
from sklearn.preprocessing import LabelEncoder
# ── Broker-accurate accounting (Factual Trading Product, Phase 1) ──────────
# These mirror v2/config.py but are duplicated HERE because strategy_utils is
# INLINED standalone into Modal training jobs (no `config` import available in
# the container). Keep in sync with config.py if those change.
LOT_UNITS = 100_000 # units per standard lot
COMMISSION_USD = 6.0 # round-trip commission, USD per standard lot
SWAP_LONG_USD = -1.00 # USD per lot per night held long (cost)
SWAP_SHORT_USD = 0.30 # USD per lot per night held short (credit)
DEFAULT_FX_SPREAD_PIPS = 0.15 # default bid/ask spread (pips) for FX majors (#283)
def _usd_conv(symbol, rate):
"""Divisor to convert quote-currency P&L into the USD account currency.
XXXUSD pairs (EUR/GBP/AUD/NZD-USD) are already quoted in USD → conv = 1.
USDXXX pairs (USD-JPY/CHF/CAD) are quoted in the foreign ccy → divide the
quote-ccy P&L by the pair's price (use the trade's exit price; brokers
convert at close). Replaces the old "EURUSD-approximated for all" caveat.
"""
s = (symbol or "").upper().replace("/", "")
if len(s) >= 6 and s[:3] == "USD":
return rate if rate else 1.0
return 1.0
# Keep in sync with predictor._CRYPTO_PIP (self-contained; strategy_utils has no config import).
_CRYPTO_PIP = {"BTCUSD": 1.0}
def _pip_size(symbol):
"""Pip size: crypto uses its tick (BTCUSD=1.0), JPY pairs 0.01, other majors 0.0001."""
s = (symbol or "").upper()
if s in _CRYPTO_PIP:
return _CRYPTO_PIP[s]
return 0.01 if s.endswith("JPY") else 0.0001
def _default_spread_pips(symbol):
"""Default bid/ask spread (pips) when the caller doesn't specify (#283).
FX majors ~0.15 pip (raw/ECN — pairs with the $6/lot commission). Crypto
returns 0.0: its costs are modeled via HL maker/taker tiers, not a spread."""
s = (symbol or "").upper()
if s in _CRYPTO_PIP:
return 0.0
return DEFAULT_FX_SPREAD_PIPS
def _trade_pnl_usd(direction, entry_px, exit_px, lots, symbol, nights,
commission_usd=None, swap_long_usd=None, swap_short_usd=None):
"""Broker-accurate per-trade P&L in USD (account currency).
units = lots × LOT_UNITS
pnl_quote = units × (exit − entry) × dir (dir +1 buy / −1 sell)
pnl_usd = pnl_quote / conv(symbol, exit_px) (per-pair FX)
pnl_usd −= commission × lots (round-trip)
pnl_usd += swap(side) × lots × nights (swap_long is negative)
commission/swap default to the module constants when None — the Config
tab / API values are threaded through run_strategy (#254: they used to be
collected but never consumed, so tuning them silently did nothing).
"""
_comm = COMMISSION_USD if commission_usd is None else float(commission_usd)
_swl = SWAP_LONG_USD if swap_long_usd is None else float(swap_long_usd)
_sws = SWAP_SHORT_USD if swap_short_usd is None else float(swap_short_usd)
units = lots * LOT_UNITS
pnl_quote = units * (exit_px - entry_px) * direction
usd = pnl_quote / _usd_conv(symbol, exit_px)
usd -= _comm * lots
usd += (_swl if direction > 0 else _sws) * lots * nights
return usd
def _pips(symbol, direction, entry_px, exit_px):
"""Signed pips for the trade (cTrader-parity Pips column)."""
return direction * (exit_px - entry_px) / _pip_size(symbol)
def _nights_held(idx_dates, idx_dow, entry_bar, exit_bar):
"""Overnight rolls between entry and exit bars. Each calendar-day boundary
the position is held across = one night; Wednesday→Thursday rolls ×3 (covers
weekend settlement) — matches the legacy pipeline.py swap model."""
if idx_dates is None or entry_bar is None or exit_bar is None:
return 0.0
nights = 0.0
for j in range(int(entry_bar), int(exit_bar)):
if idx_dates[j] != idx_dates[j + 1]:
nights += 3.0 if idx_dow[j] == 2 else 1.0
return nights
def _index_date_arrays(idx):
"""Return (dates, dayofweek) arrays for a DatetimeIndex, or (None, None)."""
try:
return idx.date, idx.dayofweek.values
except Exception:
try:
return (np.array([t.date() for t in idx]),
np.array([t.dayofweek for t in idx]))
except Exception:
return None, None
# Max backtest window per timeframe. A finer timeframe over a longer window
# blows up the results dict / parquet load / Modal train time (the 2026-05-12
# OOM was a 1-min × multi-year sweep) — and a 1-min strategy gains nothing from
# 2 years of 1-min bars. Enforced HERE because every training path (UI / API /
# Modal) funnels through run_strategy → load_ohlc. Env-overridable so a future
# "max plan" / dedicated-server tier can lift it.
_TF_MAX_DAYS = {
"1min": 30,
"5min": 90,
"15min": 365,
"1h": 730,
}
def tf_max_days(tf):
"""Public accessor for the per-timeframe window cap, env override applied.
The dash layer uses this to clamp cfg-start-date at entry (#364) so the UI
agrees with the silent load_ohlc truncation below. None = no cap known.
"""
import os as _os
if not tf:
return None
_max = _TF_MAX_DAYS.get(str(tf).strip())
if not _max:
return None
_env = _os.environ.get(f"QM_MAX_DAYS_{str(tf).strip().upper()}")
if _env and _env.isdigit():
_max = int(_env)
return _max
def clamp_backtest_start(start_date, end_date, tf):
"""(start_out, was_clamped, max_days) — pure window-cap policy (#364).
Accepts YYYY-MM-DD or DD-MM-YYYY strings; a clamped start comes back
canonical YYYY-MM-DD. Unparseable input or an uncapped timeframe passes
through untouched. Missing/unparseable end defaults to today.
"""
_max = tf_max_days(tf)
if not _max:
return start_date, False, None
def _parse(s):
for _fmt in ("%Y-%m-%d", "%d-%m-%Y", "%d-%m-%y"):
try:
return pd.Timestamp(__import__("datetime").datetime.strptime(str(s).strip(), _fmt))
except (ValueError, TypeError):
continue
return None
_start = _parse(start_date)
if _start is None:
return start_date, False, _max
_end = _parse(end_date) or pd.Timestamp.utcnow().normalize().tz_localize(None)
_floor = _end - pd.Timedelta(days=_max)
if _start < _floor:
return _floor.strftime("%Y-%m-%d"), True, _max
return start_date, False, _max
def normalize_bt_date(s):
"""'YYYY-MM-DD[ HH:MM[:SS]]' | 'DD-MM-YYYY' | 'DD-MM-YY' → 'YYYY-MM-DD', else None.
#373: bt_config carries a MIX of exact datetimes (vline drags) and bare
dates (typed input). Naive string compares between the two are prefix
accidents ('2026-07-22 23:45:00' >= '2026-07-22' is True) — every date
comparison must go through this date-part normalisation first.
"""
if s is None:
return None
_s = str(s).strip()
if not _s:
return None
# Drop a time component if present — the date part is what we compare.
_date_part = _s.split(" ")[0].split("T")[0]
import datetime as _dt
for _fmt in ("%Y-%m-%d", "%d-%m-%Y", "%d-%m-%y"):
try:
return _dt.datetime.strptime(_date_part, _fmt).strftime("%Y-%m-%d")
except ValueError:
continue
return None
def validate_date_order(start_date, end_date):
"""(ok, start_norm, end_norm) — reject start >= end at DATE granularity (#373).
ok=False only when BOTH sides parse and start_norm >= end_norm; missing or
unparseable input passes (never block on garbage — downstream defaults
handle it, same stance as clamp_backtest_start).
"""
_sn = normalize_bt_date(start_date)
_en = normalize_bt_date(end_date)
if _sn and _en and _sn >= _en:
return False, _sn, _en
return True, _sn, _en
def clamp_to_data_tail(start_date, end_date, tail_date):
"""(start_out, end_out, clamped) — snap date-parts past the data tail back (#373).
tail_date is the last available candle's 'YYYY-MM-DD'. A start or end whose
date-part lies beyond it comes back as tail_date, with 'start'/'end' added
to the `clamped` set. Unparseable/empty values pass through untouched
(the JS date poller already ignores garbage).
"""
clamped = set()
_tail = normalize_bt_date(tail_date)
if not _tail:
return start_date, end_date, clamped
_sn = normalize_bt_date(start_date)
_en = normalize_bt_date(end_date)
if _sn and _sn > _tail:
start_date, clamped = _tail, clamped | {"start"}
if _en and _en > _tail:
end_date, clamped = _tail, clamped | {"end"}
return start_date, end_date, clamped
DEFAULT_BACKTEST_DAYS = 40
def resolve_backtest_window(start_date, end_date, tf):
"""(start, end, was_clamped, max_days) — the ONE default-window policy (#370).
Every deploy surface resolves its window here so they can't drift: /landing
and the /docs Quick Deploy widget (via /api/v1/one_shot) and the /docs
generate→train→deploy composer (via /api/v1/train) all land on the same
dates. Before this, the composer carried its own 28d/-5d window hardcoded in
the page — a client-side mirror of a server policy, which is the same shape
of bug as #367.
Missing end defaults to today, so the stats cover up to the latest bar
(user 2026-05-24: "I want it till Today!"). Missing start defaults to
DEFAULT_BACKTEST_DAYS back. The result is then clamped to the timeframe's
window cap, so the requested window is always the window that trains.
"""
_end = str(end_date or "").strip()
_start = str(start_date or "").strip()
_today = pd.Timestamp.utcnow().normalize().tz_localize(None)
if not _end:
_end = _today.strftime("%Y-%m-%d")
if not _start:
_base = pd.Timestamp(_end) if _end else _today
_start = (_base - pd.Timedelta(days=DEFAULT_BACKTEST_DAYS)).strftime("%Y-%m-%d")
_start, _capped, _max = clamp_backtest_start(_start, _end, tf)
return _start, _end, _capped, _max
def _fetch_ohlc_from_internal(symbol: str, tf: str, start: str, end: str):
"""Phase 3.2: fetch parquet bytes from Server A's /internal/ohlc endpoint
instead of reading a local file. Used inside Modal containers / Mac worker
pool (Phase 3.4) so every train sees the same source of truth as the chart.
Returns: pd.DataFrame (parquet decoded), or raises on any failure so the
caller can fall back / surface a clear error in the job.
"""
import hashlib as _hashlib, hmac as _hmac, io as _io, os as _os
import urllib.request as _ur, urllib.parse as _urp
base = (_os.environ.get("QM_INTERNAL_OHLC_BASE") or "").rstrip("/")
secret = (_os.environ.get("INTERNAL_WS_SECRET") or "").strip()
if not base:
raise RuntimeError("QM_INTERNAL_OHLC_BASE not set")
if not secret:
raise RuntimeError("INTERNAL_WS_SECRET not set")
msg = f"{symbol}|{tf}|{start}|{end}".encode("utf-8")
sig = _hmac.new(secret.encode("utf-8"), msg, _hashlib.sha256).hexdigest()
qs = _urp.urlencode({
"symbol": symbol, "tf": tf,
"start": start, "end": end, "sig": sig,
})
url = f"{base}/internal/ohlc?{qs}"
req = _ur.Request(url, headers={"User-Agent": "qm-worker/1.0"})
with _ur.urlopen(req, timeout=30) as resp:
if resp.status != 200:
raise RuntimeError(f"/internal/ohlc returned {resp.status}")
payload = resp.read()
print(f"[load_ohlc:internal] {symbol} {tf} fetched {len(payload)} bytes", flush=True)
return pd.read_parquet(_io.BytesIO(payload))
def _parse_symbol_tf_from_path(data_path: str):
"""Pull SYMBOL + TF out of a path like .../EURUSD_1min.parquet."""
import os as _os, re as _re
base = _os.path.basename(str(data_path))
m = _re.match(r"^([A-Z]{6})_(\d+min|\d+h)\.parquet$", base)
if not m:
return None, None
return m.group(1), m.group(2)
def load_ohlc(data_path, start_date="", end_date=""):
"""Load OHLC parquet, sort index, filter dates. Always returns consistent format.
The lower bound is clamped per timeframe (see _TF_MAX_DAYS) — a request for
more history than the cap silently starts later.
Phase 3.2: when env QM_USE_INTERNAL_OHLC=="1", fetch over HTTP from
Server A's /internal/ohlc endpoint instead of pd.read_parquet on a local
file (which on Modal is a stale Volume snapshot). The endpoint applies the
same day-cap, so the local cap-check below is a defensive no-op in that
path. Flag defaults to "0" → unchanged behavior.
Returns: (df, close, open_, high, low)
"""
import os as _os, re as _re
_use_internal = _os.environ.get("QM_USE_INTERNAL_OHLC", "0") == "1"
if _use_internal:
_sym, _tf = _parse_symbol_tf_from_path(data_path)
if not _sym or not _tf:
raise RuntimeError(
f"QM_USE_INTERNAL_OHLC=1 but DATA_PATH basename does not match "
f"SYMBOL_TF.parquet: {data_path}"
)
df = _fetch_ohlc_from_internal(_sym, _tf, start_date or "", end_date or "")
else:
df = pd.read_parquet(data_path)
df.index = pd.to_datetime(df.index)
df = df.sort_index()
# Per-timeframe window cap (timeframe inferred from the parquet filename).
_m = _re.search(r"_(\d+min|\d+h)\.parquet$", _os.path.basename(str(data_path)))
_tf = _m.group(1) if _m else None
_max_days = tf_max_days(_tf)
if _max_days and _max_days > 0 and len(df):
try:
_eff_end = pd.Timestamp(end_date) if end_date else df.index.max()
_eff_end = min(_eff_end, df.index.max())
_floor = _eff_end - pd.Timedelta(days=_max_days)
_req_start = pd.Timestamp(start_date) if start_date else df.index.min()
if _req_start < _floor:
print(f"[load_ohlc] {_tf} backtest window capped to {_max_days}d: "
f"start {_req_start.date()} -> {_floor.date()}", flush=True)
start_date = _floor
except Exception as _e:
print(f"[load_ohlc] window-cap check skipped ({_e})", flush=True)
if start_date:
df = df[df.index >= start_date]
if end_date:
df = df[df.index <= end_date]
return df, df["close"], df["open"], df["high"], df["low"]
TF_MINUTES = {"1min": 1, "5min": 5, "15min": 15, "1h": 60}
MAX_AUX_FEEDS = 3
AUX_MIN_COVERAGE = 0.90 # completed aux bar must exist for >=90% of primary
# rows *after* the feed has warmed up (see note below)
def normalize_aux_specs(raw, primary_symbol, primary_tf):
"""#387/#388: validate/normalize aux feed specs. Returns (specs, warnings).
Honest-drop discipline (#377): every removal produces a user-facing
warning that is true in every branch - never silently shrink."""
from data_paths import SUPPORTED_SYMBOLS, SUPPORTED_TIMEFRAMES
specs, warnings, seen = [], [], set()
for item in (raw or []):
if not isinstance(item, dict):
warnings.append("an aux feed entry was not an object - ignored.")
continue
sym = str(item.get("symbol") or primary_symbol).upper().strip()
tf = str(item.get("timeframe") or primary_tf).strip()
if sym not in SUPPORTED_SYMBOLS:
warnings.append(f"aux feed {sym} {tf}: unsupported symbol - dropped.")
continue
if tf not in SUPPORTED_TIMEFRAMES:
warnings.append(f"aux feed {sym} {tf}: unsupported timeframe - dropped.")
continue
if sym == str(primary_symbol).upper() and tf == primary_tf:
warnings.append(
f"aux feed {sym} {tf} equals the primary feed - dropped.")
continue
if (sym, tf) in seen:
continue
if len(specs) >= MAX_AUX_FEEDS:
warnings.append(
f"aux feed {sym} {tf}: over the {MAX_AUX_FEEDS}-feed cap - dropped.")
continue
seen.add((sym, tf))
specs.append({"symbol": sym, "timeframe": tf})
return specs, warnings
def _parse_aux_feeds_regex_fallback(code):
"""Line-anchored regex fallback for `parse_aux_feeds_from_code`, used
only when `ast.parse` itself fails (pasted code mid-edit and genuinely
unparsable). Deliberately limited to a single physical line ending in
`]` - it cannot see multi-line literals or trailing comments, but on
unparsable input there is no AST to walk anyway, and a still-legal
single-line stamp is worth recovering rather than giving up entirely.
Still takes the LAST match, never the first."""
import ast as _ast
import re as _re
matches = _re.findall(r"^AUX_FEEDS\s*=\s*(\[.*\])\s*$", code or "",
_re.MULTILINE)
if not matches:
return []
try:
parsed = _ast.literal_eval(matches[-1])
except Exception:
return []
return parsed if isinstance(parsed, list) else []
def parse_aux_feeds_from_code(code):
"""Task 9 (#387/#388): read the AUX_FEEDS the deployed strategy code
actually carries, for the deploy bundle / live predictor.
MUST take the LAST module-scope match, never the first. Task 7's
builder guarantees AUX_FEEDS by APPENDING its authoritative assignment
as the last module-scope statement (module execution is last-write-
wins) - the model's own earlier AUX_FEEDS line can legitimately survive
above it, because the builder's cleanliness strip deliberately declines
to remove statements it does not fully own. A first-match reader (e.g.
re.search) would read the MODEL's value - the exact exploit Task 7
spent five review rounds closing, reintroduced here.
Parses via `ast` rather than a line-anchored regex: walks the module's
top-level statements for `Assign` nodes with a single, bare
`Name("AUX_FEEDS")` target, keeps the LAST one, and `literal_eval`s its
value. This handles a multi-line literal and a trailing same-line
comment correctly (both legal Python the BACKTEST already handles fine
by executing the module - only this static reader used to choke on
them), and it structurally cannot mistake an "AUX_FEEDS = [...]" string
that only appears inside a docstring for a real assignment (a
docstring's TEXT is never re-parsed as further statements).
Falls back to a line-anchored regex only when `ast.parse` itself fails
(pasted code can be genuinely mid-edit / unparsable) - the regex path
still recovers a still-legal single-line stamp elsewhere in text that
doesn't compile, rather than giving up outright.
Returns [] on any failure (missing, malformed, or not a list) - never
raises."""
import ast as _ast
code = code or ""
try:
tree = _ast.parse(code)
except Exception:
return _parse_aux_feeds_regex_fallback(code)
last_value_node = None
for node in tree.body:
if (isinstance(node, _ast.Assign)
and len(node.targets) == 1
and isinstance(node.targets[0], _ast.Name)
and node.targets[0].id == "AUX_FEEDS"):
last_value_node = node.value
if last_value_node is None:
return []
try:
parsed = _ast.literal_eval(last_value_node)
except Exception:
return []
return parsed if isinstance(parsed, list) else []
def _default_aux_loader(symbol, tf, start=None, end=None):
"""#403: aux data must travel the SAME road as the primary feed.
load_ohlc() switches to /internal/ohlc when QM_USE_INTERNAL_OHLC=1, which
is exactly what Modal containers set - and since M22c removed the Volume
there is no OHLC on the container filesystem at all. A loader that reads
a local parquet can therefore only ever raise there, and merge_aux_feeds
turns that into a silent "dropped" - the model trains with no aux columns
and still reports success. One transport decision, made in one place.
"""
import os as _os
if _os.environ.get("QM_USE_INTERNAL_OHLC", "0") == "1":
return _fetch_ohlc_from_internal(symbol, tf, start or "", end or "")
from data_paths import data_path_for
aux = pd.read_parquet(data_path_for(symbol, tf))
# #406: honour the window in BOTH branches. The HTTP branch gets it for
# free (the endpoint filters), so leaving the local branch unfiltered
# made the same call mean different things depending on an env var - and
# the finer-than-primary merge below then ran three groupby.apply passes
# over the ENTIRE 1min history instead of the few days actually being
# merged. Measured on Server A: a 1min feed into a 5min primary took
# 128.78s unfiltered vs 0.04s for an equal-timeframe feed.
if start or end:
aux = aux.copy()
aux.index = pd.to_datetime(aux.index)
if start:
aux = aux[aux.index >= pd.Timestamp(start)]
if end:
aux = aux[aux.index <= pd.Timestamp(end)]
return aux
def merge_aux_feeds(df, aux_specs, primary_symbol, primary_tf, loader=None):
"""#387/#388: merge aux (symbol, timeframe) feeds into the primary df as
aux_* columns, BEFORE feature_engineering. The single audited home of the
no-lookahead rule - backtest (run_strategy) and live (predictor) both call
this, so train/serve skew is structurally impossible.
Info-time of primary row t is t + primary_tf (the bar's own close is known
then - identical to how the existing pipeline treats the bar's own OHLC).
An aux bar labeled s (timeframe atf) becomes visible at s + atf.
Loader contract (#403): loader(symbol, timeframe, start, end) -> DataFrame.
start/end are ISO strings bracketing the window this merge needs; a live
loader that always serves the latest candles may ignore them.
"""
loader = loader or _default_aux_loader
report = {"merged": [], "dropped": []}
# CC-2 (final whole-branch review): TF_MINUTES is a SECOND source of
# truth for supported timeframes, separate from SUPPORTED_TIMEFRAMES /
# normalize_aux_specs' own validation - they can drift apart (e.g. a
# future SUPPORTED_TIMEFRAMES addition this dict never learns about).
# run_strategy calls this function for EVERY strategy, aux or not, so a
# primary_tf missing here must never raise: the NO-AUX path (the vast
# majority of strategies) must return the df untouched exactly like the
# `not aux_specs` early-out below, and the WITH-AUX path must drop every
# requested feed with an honest reason naming the unsupported primary
# timeframe - never a bare KeyError either way.
ptf_min = TF_MINUTES.get(primary_tf)
if ptf_min is None:
if not aux_specs:
return df, report
for spec in aux_specs:
report["dropped"].append({
"symbol": spec.get("symbol"), "timeframe": spec.get("timeframe"),
"reason": f"primary timeframe {primary_tf!r} is not supported "
f"for aux merging"})
return df, report
if not aux_specs:
return df, report
# pd.merge_asof silently produces WRONG results (not an error) when its
# LEFT frame isn't sorted by the join key. `aux` is sorted defensively
# below but df's own index wasn't - sort it too.
df = df.sort_index()
info_time = df.index + pd.Timedelta(minutes=ptf_min)
for spec in aux_specs:
sym, atf = spec["symbol"], spec["timeframe"]
atf_min = TF_MINUTES[atf]
prefix = f"aux_{sym.lower()}_{atf}"
# #403: ask for the window we actually need. /internal/ohlc day-caps
# an unbounded request to the most RECENT window, so a historical
# backtest that asked for "everything" would get aux data that does
# not overlap the primary at all - and the overlap failure is silent
# (all-NaN columns, then a coverage drop). The pad covers an FX
# weekend plus one aux bar so row 0 already has a completed aux bar.
_pad = pd.Timedelta(days=3) + pd.Timedelta(minutes=atf_min)
_start = (df.index.min() - _pad).strftime("%Y-%m-%d %H:%M:%S")
_end = (df.index.max() + pd.Timedelta(minutes=ptf_min)).strftime(
"%Y-%m-%d %H:%M:%S")
try:
aux = loader(sym, atf, _start, _end)
except Exception as exc:
report["dropped"].append({"symbol": sym, "timeframe": atf,
"reason": f"load failed: {exc}"})
continue
if aux is None or len(aux) == 0:
report["dropped"].append({"symbol": sym, "timeframe": atf,
"reason": "no data"})
continue
aux = aux.sort_index()
aux.index = pd.to_datetime(aux.index)
cols = []
if atf_min >= ptf_min:
# Completed-bar view: asof on COMPLETION time (label + atf).
right = aux[["open", "high", "low", "close"]].copy()
right.index = right.index + pd.Timedelta(minutes=atf_min)
right.columns = [f"{prefix}_{c}" for c in right.columns]
merged = pd.merge_asof(
pd.DataFrame(index=info_time), right,
left_index=True, right_index=True, direction="backward")
merged.index = df.index # back to primary labels
# Coverage is measured only from the first row where a completed
# aux bar is even possible - every coarser/equal feed has an
# unavoidable, correct, all-NaN warm-up before its first bar
# completes (see test_rows_before_first_aux_bar_are_nan); that
# is not a data-quality problem and must not count against the
# feed. Gaps AFTER warm-up (missing/short aux history, holes in
# the source data) still count and can still trip the floor.
close_col = merged[f"{prefix}_close"]
first_valid = close_col.first_valid_index()
if first_valid is None:
coverage = 0.0
else:
coverage = close_col.loc[first_valid:].notna().mean()
if coverage < AUX_MIN_COVERAGE:
report["dropped"].append({
"symbol": sym, "timeframe": atf,
"reason": f"completed-bar coverage {coverage:.0%} < "
f"{AUX_MIN_COVERAGE:.0%} of primary rows"})
continue
for c in merged.columns:
df[c] = merged[c]
cols.append(c)
if atf_min > ptf_min:
# Forming-bar view (Malcolm's "the 1h bar up till 11:15"):
# a running OHLC aggregate over the aux symbol's PRIMARY-tf
# bars inside the CURRENT coarse window, truncated at each
# row's own info-time. Leak-free by construction - only bars
# that have already closed enter the running aggregate.
# Same-symbol feeds reuse df's own OHLC (no extra load);
# cross-symbol feeds load the aux symbol at the primary tf.
try:
fine = (df if sym == str(primary_symbol).upper()
else loader(sym, primary_tf))
fine = fine.sort_index()
fine.index = pd.to_datetime(fine.index)
win = fine.index.floor(f"{atf_min}min")
g = fine.groupby(win)
forming = pd.DataFrame({
f"{prefix}_live_open": g["open"].transform("first"),
f"{prefix}_live_high": g["high"].cummax(),
f"{prefix}_live_low": g["low"].cummin(),
f"{prefix}_live_close": fine["close"],
"_window": win,
}, index=fine.index)
# visible once THIS fine bar itself completes
forming.index = forming.index + pd.Timedelta(minutes=ptf_min)
fm = pd.merge_asof(
pd.DataFrame(index=info_time), forming,
left_index=True, right_index=True, direction="backward")
fm.index = df.index
# A forming value carried over from a PREVIOUS coarse
# window (e.g. a gap in the fine data) must not leak into
# a new window before that window's own first fine bar
# has closed - mask those rows back to NaN.
row_window = df.index.floor(f"{atf_min}min")
live_cols = [c for c in fm.columns if c != "_window"]
stale = fm["_window"].values != row_window.values
fm.loc[stale, live_cols] = np.nan
for c in live_cols:
df[c] = fm[c]
cols.append(c)
except Exception as exc:
print(f"[aux_feeds] forming-bar columns for {sym} {atf} "
f"skipped ({exc})", flush=True)
else:
# Finer-than-primary: raw fine OHLC aggregates straight back into
# the primary bar (adds nothing at bar-close granularity), so
# provide within-bar PATH features instead. Fine bars labeled in
# [t, t+ptf) all complete by row t's info-time (t+ptf) - no shift
# needed, unlike the coarser branch above. `floor` groups each
# fine bar into the primary bar that CONTAINS it (same alignment
# assumption the forming-bar branch already makes) - a fine bar
# from the NEXT primary bar must never land in this group, or
# that's lookahead (see test_finer_uses_only_own_bar).
window = aux.index.floor(f"{ptf_min}min")
g = aux.groupby(window)
# pct_change/cummax must reset at each group's own start (a
# global pct_change would leak the previous bar's last close
# into this bar's first fine-bar return).
#
# #407: these were three groupby.apply calls carrying Python
# lambdas - one interpreter round-trip PER PRIMARY BAR. Cost
# tracked group COUNT, not row count: 200k fine bars took 14.63s
# against a 5min primary (40k groups) but 4.55s against a 15min
# one (13.3k groups) - same data, 3x the groups, 3.2x the time.
# The groupby transforms below reset per group exactly as the
# lambdas did, so the numbers are unchanged (test_407 pins them
# against the original implementation as an oracle).
_close = aux["close"]
_cg = _close.groupby(window)
rvol = _cg.pct_change().groupby(window).std()
_cummax = _cg.cummax()
maxdd = ((_cummax - _close) / _cummax).groupby(window).max()
upratio = (aux["close"] > aux["open"]).groupby(window).mean()
feats = pd.DataFrame({
f"{prefix}_rvol": rvol,
f"{prefix}_upratio": upratio,
f"{prefix}_maxdd": maxdd,
})
feats = feats.reindex(df.index)
# Coverage measures "did we have ANY fine data for this primary
# bar" via group non-emptiness (g.size(), reindexed) - NOT via
# the `_rvol` column. `rvol` needs >=3 fine bars in a group to be
# non-NaN (2 to get a single pct-change, 2 valid values to get a
# std) while `upratio`/`maxdd` only need 1 - keying coverage to
# `rvol` would drop an otherwise-good sparse feed (e.g. exactly
# 2 fine bars/primary bar) entirely, even though 2 of 3 columns
# are fully populated. As with the completed-bar branch above,
# measured from the first row a group is even possible onward -
# an unavoidable leading warm-up before the fine feed's first
# group is not a data-quality problem and must not count
# against the feed.
group_sizes = g.size().reindex(df.index)
first_valid = group_sizes.first_valid_index()
if first_valid is None:
coverage = 0.0
else:
coverage = group_sizes.loc[first_valid:].notna().mean()
if coverage < AUX_MIN_COVERAGE:
report["dropped"].append({
"symbol": sym, "timeframe": atf,
"reason": f"fine-bar coverage {coverage:.0%} < "
f"{AUX_MIN_COVERAGE:.0%} of primary rows"})
continue
for c in feats.columns:
df[c] = feats[c]
cols.append(c)
report["merged"].append({"symbol": sym, "timeframe": atf,
"columns": cols})
return df, report
# Suffix -> meaning, checked longest-first so "aux_x_1h_live_close" resolves
# to the live_* meaning instead of the plain "close" one (both end in
# "_close" - order is what disambiguates them, not a smarter match).
_AUX_SUFFIX_ORDER = (
"live_open", "live_high", "live_low", "live_close",
"open", "high", "low", "close",
"rvol", "upratio", "maxdd",
)
def _aux_column_meaning(col, sym, tf):
"""#387/#388: one-line, per-suffix meaning for an aux_* column, used to
teach codegen what already exists on df so it never recomputes it."""
for suf in _AUX_SUFFIX_ORDER:
if col.endswith("_" + suf):
if suf.startswith("live_"):
return (f"the FORMING {tf} bar of {sym}, truncated at the "
f"current bar - no lookahead")
if suf in ("open", "high", "low", "close"):
return f"the last COMPLETED {tf} bar of {sym} ({suf})"
if suf == "rvol":
return f"std of {tf} returns within the current bar"
if suf == "upratio":
return f"fraction of up {tf} bars within the current bar"
if suf == "maxdd":
return "worst drawdown within the current bar"
# Defensive fallback - every column merge_aux_feeds produces matches one
# of the suffixes above; this only fires if that ever drifts.
return f"an aux feed column from {sym} {tf}"
def aux_columns_prompt(report):
"""#387/#388: turn a merge_aux_feeds() report into a codegen prompt block
teaching the LLM which aux_* columns already sit on df, so it never
recomputes them and never assigns to them. "" when nothing merged - the
no-aux generation path must gain ZERO content from this function."""
merged = (report or {}).get("merged") or []
if not merged:
return ""
lines = [
"",
"",
"AUX FEED COLUMNS (already merged onto df before feature_engineering runs):",
]
for entry in merged:
sym = entry.get("symbol", "")
tf = entry.get("timeframe", "")
for col in entry.get("columns") or []:
lines.append(f" - {col}: {_aux_column_meaning(col, sym, tf)}")
lines.append(
"These columns already exist on df - never recompute them and never "
"read any other data source.")
lines.append("Never assign to any aux_* column.")
return "\n".join(lines)
def make_target(close, horizon=4):
"""Create target: direction N bars ahead. Default 4 bars = 1 hour on 15-min data.
Returns: target (pd.Series of -1, 0, 1)
"""
return np.sign(close.shift(-horizon) - close)
def split_bounds(n_rows, split_idx, purge_bars=0, embargo_bars=0):
"""#411: the ONE chronological splitter. Train slice is [:train_end],
test slice is [test_start:].
purge_bars: rows removed from the END of train. make_target looks
`target_horizon` bars forward, so without purging the last `horizon`
train rows carry labels computed from the first test rows - leakage
across the boundary. purge_bars=target_horizon closes it.
embargo_bars: extra gap before test starts (serial-correlation guard).
Defaults (0, 0) reproduce the legacy split EXACTLY - the headline 70/30
path must stay bit-identical (see TestHeadlineUnchanged).
"""
split_idx = max(1, min(int(split_idx), n_rows - 1))
train_end = max(1, split_idx - int(purge_bars))
test_start = min(n_rows - 1, split_idx + int(embargo_bars))
if train_end > test_start:
train_end = test_start
return train_end, test_start
def split_data(df, target, feature_cols, train_split=0.7, validation_date=""):
"""Train/test split. Handles both ratio and date-based splits.
Drops NaN from target before splitting. Encodes labels to [0,1,2].
Returns: dict with keys:
X_train, X_test, y_train, y_test,
y_train_enc, y_test_enc, enc,
close_train, close_test,
split_idx, split_dt, n_train, n_test
"""
# Drop NaN from target
mask = target.notna()
df = df[mask].copy()
target = target[mask]
close = df["close"]
# Build feature matrix
X = df[feature_cols].copy()
X = X.bfill().ffill()
X = X.replace([np.inf, -np.inf], np.nan).fillna(0.0)
# Split
if validation_date:
split_idx = len(df[df.index <= validation_date])
else:
split_idx = int(len(df) * train_split)
# #411 - route through the single chronological splitter (default
# purge/embargo = 0, so this is bit-identical to the old inline math).
split_idx, _ = split_bounds(len(df), split_idx)
X_train = X.iloc[:split_idx]
X_test = X.iloc[split_idx:]
y_train = target.iloc[:split_idx]
y_test = target.iloc[split_idx:]
close_train = close.iloc[:split_idx]
close_test = close.iloc[split_idx:]
split_dt = str(df.index[split_idx])
# Label encoding — always fit on [-1, 0, 1]
enc = LabelEncoder()
enc.fit([-1, 0, 1])
y_train_enc = enc.transform(y_train)
y_test_enc = enc.transform(y_test)
return {
"df": df, "X_train": X_train, "X_test": X_test,
"y_train": y_train, "y_test": y_test,
"y_train_enc": y_train_enc, "y_test_enc": y_test_enc,
"enc": enc,
"close": close, "close_train": close_train, "close_test": close_test,
"split_idx": split_idx, "split_dt": split_dt,
"n_train": len(X_train), "n_test": len(X_test),
}
def compute_overlays(close, df_index):
"""Compute BB and MA overlays on full dataset. Always consistent.
Returns: (bb_dict, ma_dict)
"""
bb_mid = close.rolling(20).mean()
bb_std = close.rolling(20).std()
bb_upper = bb_mid + 2 * bb_std
bb_lower = bb_mid - 2 * bb_std
ma50 = close.rolling(50).mean()
ma100 = close.rolling(100).mean()
ma200 = close.rolling(200).mean()
def _safe(s):
s = s.reindex(df_index).bfill().ffill()
return [float(x) if (x is not None and not np.isnan(x) and not np.isinf(x)) else None
for x in s.values]
bb = {"upper": _safe(bb_upper), "mid": _safe(bb_mid), "lower": _safe(bb_lower)}
ma = {"ma50": _safe(ma50), "ma100": _safe(ma100), "ma200": _safe(ma200)}
return bb, ma
def run_backtest(signal, close, capital=10000, cost=2e-5, leverage=1.0,
lot_size=None, balance=None, symbol="",
commission_usd=None, swap_long_usd=None, swap_short_usd=None,
spread_pips=0.0):
"""Run backtest with transaction costs.
Uses price-based trade returns (same as webapp _compute_trades).
Signal 0 = hold (keep current position), not close.
Two accounting models:
• `lot_size is None` → LEGACY: per-trade return = leverage × fractional
return-on-equity (clamp −100%); `pnl` is the raw price delta; no `$`.
Direct callers (e.g. live_track / Server G) stay on this path.
• `lot_size` set → BROKER-ACCURATE (Phase 1): per-trade `pnl_usd` via
position size × per-pair FX − commission/swap; `pnl_pct` becomes
return-on-equity = pnl_usd / balance (leverage no longer multiplies).
Returns: dict with equity, trade_returns, long_returns, short_returns, bar_returns
"""
def _lev(r):
return max(r * leverage, -1.0)
broker = lot_size is not None
bal = float(balance) if balance is not None else float(capital)
if broker and bal <= 0:
bal = float(capital) or 10000.0
sig_arr = signal.values
price_arr = close.values
idx = signal.index
n = len(price_arr)
_idx_dates, _idx_dow = _index_date_arrays(idx) if broker else (None, None)
_spread_px = float(spread_pips or 0.0) * _pip_size(symbol) # bid/ask cost (#283)
def _close_ret(direction, e_px, x_px, e_bar, x_bar):
"""Return (ret_for_equity, pnl_usd_or_None) for a closed trade."""
if broker:
nights = _nights_held(_idx_dates, _idx_dow, e_bar, x_bar)
usd = _trade_pnl_usd(direction, e_px, x_px, lot_size, symbol, nights,
commission_usd=commission_usd,
swap_long_usd=swap_long_usd,
swap_short_usd=swap_short_usd)
# one full round-trip spread (buy@ask + sell@bid), in USD (#283)
usd -= _spread_px * (lot_size or 0.0) * LOT_UNITS / _usd_conv(symbol, x_px)
return max(usd / bal, -1.0), usd
# legacy leverage mode: subtract spread as a fraction of entry price (#283)
return _lev(float(direction * (x_px - e_px) / e_px - cost - _spread_px / e_px)), None
# Trade returns — price-based (matches webapp _compute_trades exactly)
trade_returns = []
long_returns = []
short_returns = []
trade_log = []
last_dir = None
entry_price = None
entry_bar = None
for i in range(n):
s = sig_arr[i]
c = price_arr[i]
if s != 0.0 and s != last_dir:
# Direction change — close previous trade, open new
if last_dir is not None and entry_price is not None and entry_price != 0:
ret, _usd = _close_ret(last_dir, entry_price, c, entry_bar, i)
trade_returns.append(ret)
if last_dir == 1:
long_returns.append(ret)
else:
short_returns.append(ret)
trade_log.append({
"type": "Buy" if last_dir == 1 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[i]),
"entry_price": round(entry_price, 5),
"exit_price": round(c, 5),
"pnl": round(last_dir * (c - entry_price), 5),
"pnl_usd": round(_usd, 2) if _usd is not None else None,
"pips": round(_pips(symbol, last_dir, entry_price, c), 1) if broker else None,
"pnl_pct": round(ret * 100, 3),
"exit_reason": "signal",
})
entry_price = c
entry_bar = i
last_dir = s
# Close last open trade
if last_dir is not None and entry_price is not None and n > 0 and entry_price != 0:
c = price_arr[-1]
ret, _usd = _close_ret(last_dir, entry_price, c, entry_bar, n - 1)
trade_returns.append(ret)
if last_dir == 1:
long_returns.append(ret)
else:
short_returns.append(ret)
trade_log.append({
"type": "Buy" if last_dir == 1 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[-1]),
"entry_price": round(entry_price, 5),
"exit_price": round(c, 5),
"pnl": round(last_dir * (c - entry_price), 5),
"pnl_usd": round(_usd, 2) if _usd is not None else None,
"pips": round(_pips(symbol, last_dir, entry_price, c), 1) if broker else None,
"pnl_pct": round(ret * 100, 3),
"exit_reason": "end",
})
# Equity curve from trade returns
cumret = 1.0
equity_vals = np.full(n, bal if broker else float(capital))
trade_idx = 0
in_trade = False
t_entry_price = None
t_entry_bar = None
t_dir = None
for i in range(n):
s = sig_arr[i]
c = price_arr[i]
if s != 0.0 and s != t_dir:
if t_dir is not None and t_entry_price is not None and t_entry_price != 0:
t_ret, _ = _close_ret(t_dir, t_entry_price, c, t_entry_bar, i)
cumret *= (1 + t_ret)
t_entry_price = c
t_entry_bar = i
t_dir = s
equity_vals[i] = (bal if broker else capital) * cumret
# Bar returns for Sharpe
bar_returns = np.zeros(n)
for i in range(1, n):
if price_arr[i - 1] != 0 and last_dir is not None:
bar_returns[i] = sig_arr[i - 1] * (price_arr[i] - price_arr[i - 1]) / price_arr[i - 1] if sig_arr[i - 1] != 0 else 0.0
return {
"equity": pd.Series(equity_vals, index=close.index),
"trade_returns": trade_returns,
"long_returns": long_returns,
"short_returns": short_returns,
"bar_returns": bar_returns,
"trade_log": trade_log,
}
def compute_trade_stats(trades, capital=10000):
"""Single source of truth for trade statistics.
Every display path reads from this — no recomputation anywhere.
All values are rounded and JSON-safe (no inf/nan).
"""
if not trades:
return {"n": 0, "wins": 0, "losses": 0, "wr": 0, "avg": 0,
"best": 0, "worst": 0, "ret": 0, "np": 0, "mdd": 0,
"pf": 0, "rr": 0, "expect": 0}
w = [r for r in trades if r > 0]
l = [r for r in trades if r < 0]
cumret = 1.0
for r in trades:
cumret *= (1 + r)
net_p = capital * (cumret - 1)
# Max drawdown
eq = np.cumprod([1.0] + [1 + r for r in trades])
peak = np.maximum.accumulate(eq)
mdd = float(((eq - peak) / peak).min()) if len(eq) > 1 else 0.0
# Profit Factor
gross_w = sum(w) if w else 0
gross_l = abs(sum(l)) if l else 0
pf = gross_w / gross_l if gross_l > 0 else (9999.0 if gross_w > 0 else 0)
# Risk:Reward
avg_w = float(np.mean(w)) if w else 0
avg_l = abs(float(np.mean(l))) if l else 0
rr = avg_w / avg_l if avg_l > 0 else (9999.0 if avg_w > 0 else 0)
# Expectancy
expect = net_p / len(trades)
return {
"n": len(trades), "wins": len(w), "losses": len(l),
"wr": round(len(w) / len(trades), 4),
"avg": round(float(np.mean(trades)), 6),
"best": round(max(w), 6) if w else 0,
"worst": round(min(l), 6) if l else 0,
"ret": round(cumret - 1, 6),
"np": round(net_p, 2),
"mdd": round(mdd, 6),
"pf": round(pf, 2),
"rr": round(rr, 2),
"expect": round(expect, 2),
}
def compute_metrics(bt_result, close_test, capital=10000):
"""Compute all standard metrics from backtest result.
Uses trade-level compounding (same as webapp _trade_stats) for accuracy.
Returns: dict with total_ret, bh_ret, sharpe_strat, sharpe_bh, mdd, n_trades
"""
equity = bt_result["equity"]
trade_returns = bt_result["trade_returns"]
# Total return — trade-level compounding (matches webapp)
if trade_returns:
cumret = 1.0
for r in trade_returns:
cumret *= (1 + r)
total_ret = cumret - 1
else:
total_ret = 0.0
# Buy and hold
bh_equity = capital * (close_test / close_test.iloc[0])
bh_ret = (bh_equity.iloc[-1] - capital) / capital if capital != 0 else 0.0
# Sharpe ratio — trade-level (matches webapp: sqrt(252*26) annualization)
if len(trade_returns) >= 2 and float(np.std(trade_returns)) > 0:
sharpe_strat = float(np.mean(trade_returns) / np.std(trade_returns) * np.sqrt(252 * 26))
else:
sharpe_strat = 0.0
bh_rets = bh_equity.pct_change().dropna()
if len(bh_rets) > 1 and bh_rets.std() != 0:
sharpe_bh = float((bh_rets.mean() / bh_rets.std()) * np.sqrt(252 * 24 * 4))
else:
sharpe_bh = 0.0
# Max drawdown — trade-level (matches webapp)
if trade_returns:
eq = np.cumprod([1.0] + [1 + r for r in trade_returns])
peak = np.maximum.accumulate(eq)
mdd = float(((eq - peak) / peak).min()) if len(eq) > 1 else 0.0
else:
mdd = 0.0
return {
"total_ret": float(total_ret),
"bh_ret": float(bh_ret),
"sharpe_strat": float(sharpe_strat) if not np.isnan(sharpe_strat) else 0.0,
"sharpe_bh": float(sharpe_bh) if not np.isnan(sharpe_bh) else 0.0,
"mdd": float(mdd),
"n_trades": len(trade_returns),
}
# Diagnostics line/histogram series (equity / drawdown / rolling_acc / conf_hist)
# only feed the small Diagnostics charts — they're never used by the price chart
# or scroll-back. On a 1-min model trained over the (2.2-capped) window these are
# still ~30k points each; downsample to a visually-identical resolution before the
# dict leaves the trainer so it doesn't carry that into Server-A RAM / Postgres.
_RESULTS_SERIES_MAX = 5000
def _downsample_idx(n, cap=_RESULTS_SERIES_MAX):
"""Evenly-spaced index list spanning [0, n-1] (first+last always kept), or
None when no downsampling is needed (n <= cap)."""
if n <= cap:
return None
return np.unique(np.linspace(0, n - 1, cap).astype(int)).tolist()
def _take(arr, idx):
"""Subset a list by an index list (idx may be None → return arr unchanged)."""
if idx is None or not isinstance(arr, list):
return arr
return [arr[i] for i in idx]
# trade_log / train_trade_log are lists of per-trade dicts (display-only — the
# Trade Log tab). They scale with TRADE count, not bar count, so the bar-window
# cap (Phase 2.2) doesn't bound them — a degenerate near-every-bar model can put
# 10k+ trade dicts in the blob (>3 MB). Cap each (independently — a small-N model
# keeps every trade) to the most-recent N, recording `*_total` + `*_truncated`
# so the true count is still reported. Real strategies have far fewer than
# _TRADE_LOG_MAX trades, so this only ever bites pathological models.
_TRADE_LOG_MAX = 5000
def _cap_trade_log(tl):
"""Return (capped_list, original_len, was_truncated)."""
if not isinstance(tl, list) or len(tl) <= _TRADE_LOG_MAX:
return tl, (len(tl) if isinstance(tl, list) else 0), False
return tl[-_TRADE_LOG_MAX:], len(tl), True
def _safe_fi_pairs(model, feature_cols, top=15):
"""Feature-importance pairs that never crash (FACT-F2 / #3).
A non-tree model wrapped by a stale/diverged ModelWrapper can hand back
``None`` (LogReg has ``coef_``, not ``feature_importances_``). Guard here so
build_return_dict is safe regardless of which wrapper produced the PKL: fall
back to ``coef_``, then to zero-weight pairs aligned to ``feature_cols``.
"""
importances = getattr(model, "feature_importances_", None)
if importances is None:
coef = getattr(model, "coef_", None)
if coef is not None:
importances = np.abs(np.asarray(coef)).mean(axis=0)
if importances is None or len(importances) != len(feature_cols):
# length mismatch or nothing usable -> zero weights, still renders
return [(f, 0.0) for f in feature_cols][-top:]
return sorted(zip(feature_cols, importances), key=lambda x: x[1])[-top:]
def build_return_dict(split_result, bt_result, metrics, model, feature_cols,
signal_full, p_pos_test, p_neg_test, custom_figs=None,
bt_train_result=None, pre_stats=None):
"""Assemble the complete return dict. Handles ALL serialization.
Never returns Timestamps, numpy arrays, or non-JSON types.
Returns: JSON-safe dict with all required keys
"""
df = split_result["df"]
close = split_result["close"]
close_test = split_result["close_test"]
X_test = split_result["X_test"]
y_test = split_result["y_test"]
equity = bt_result["equity"]
bar_returns = bt_result["bar_returns"]
# OHLC
ohlc_dates = [str(x) for x in df.index.tolist()]
def _safe_list(arr):
return [float(x) if (x is not None and not np.isnan(x) and not np.isinf(x)) else None
for x in arr]
# Overlays
bb, ma = compute_overlays(close, df.index)
# Buy and hold equity
capital = equity.iloc[0] if len(equity) > 0 else 10000
bh_equity = capital * (close_test / close_test.iloc[0])
# Confusion matrix
from sklearn.metrics import confusion_matrix
pred_test = model.predict(X_test)
y_test_arr = np.asarray(y_test)
cm = confusion_matrix(y_test_arr, pred_test, labels=[-1, 0, 1])
# Rolling accuracy
sig_arr = signal_full.reindex(close_test.index).values
correct = pd.Series((pred_test == y_test_arr).astype(float), index=X_test.index)
active_test = pd.Series(sig_arr != 0, index=close_test.index) if len(sig_arr) == len(close_test) else pd.Series(True, index=close_test.index)
correct_active = correct.where(active_test, other=np.nan)
rolling_acc = correct_active.rolling(30, min_periods=1).mean()
# Feature importance (FACT-F2 guard: never crash on None/short importances)
fi_pairs = _safe_fi_pairs(model, feature_cols)
# Drawdown
rolling_max = equity.cummax()
drawdown = (equity - rolling_max) / rolling_max.replace(0, np.nan)
drawdown = drawdown.fillna(0.0)
# ── Downsample the Diagnostics-only series (see _downsample_idx) ──────────
_eq_dates = [str(x) for x in close_test.index.tolist()]
_eq_strat = _safe_list(equity.values)
_eq_bh = _safe_list(bh_equity.values)
_eq_idx = _downsample_idx(len(_eq_dates))
_eq_dates, _eq_strat, _eq_bh = _take(_eq_dates, _eq_idx), _take(_eq_strat, _eq_idx), _take(_eq_bh, _eq_idx)
_ra_dates = [str(x) for x in rolling_acc.index.tolist()]
_ra_vals = [float(x) if (not np.isnan(x) and not np.isinf(x)) else None for x in rolling_acc.values]
_ra_idx = _downsample_idx(len(_ra_dates))
_ra_dates, _ra_vals = _take(_ra_dates, _ra_idx), _take(_ra_vals, _ra_idx)
_dd_dates = [str(x) for x in drawdown.index.tolist()]
_dd_vals = _safe_list(drawdown.values)
_dd_idx = _downsample_idx(len(_dd_dates))
_dd_dates, _dd_vals = _take(_dd_dates, _dd_idx), _take(_dd_vals, _dd_idx)
_cp_pos = [float(x) for x in (p_pos_test.tolist() if hasattr(p_pos_test, 'tolist') else list(p_pos_test))]
_cp_neg = [float(x) for x in (p_neg_test.tolist() if hasattr(p_neg_test, 'tolist') else list(p_neg_test))]
_cp_pos = _take(_cp_pos, _downsample_idx(len(_cp_pos)))
_cp_neg = _take(_cp_neg, _downsample_idx(len(_cp_neg)))
# ── Trade logs — display-only (Trade Log tab); cap to most-recent N with a
# `_total` field so the true count is still reported (see _cap_trade_log).
# NB: ret_dist arrays are left FULL — a downstream path in callbacks.py
# recomputes n_trades/win-rate from len(ret_dist), so a sample would skew
# the displayed counts; they're small anyway and gzip handles them.
_tl_test, _tl_test_n, _tl_test_tr = _cap_trade_log(bt_result.get("trade_log", []))
_tl_tr, _tl_tr_n, _tl_tr_tr = _cap_trade_log(bt_train_result.get("trade_log", []) if bt_train_result else [])
return {
"ohlc": {
"dates": ohlc_dates,
"open": _safe_list(df["open"].values),
"high": _safe_list(df["high"].values),
"low": _safe_list(df["low"].values),
"close": _safe_list(df["close"].values),
},
"signals": {
"dates": [str(x) for x in signal_full.index.tolist()],
"values": [float(x) for x in signal_full.values],
},
"bb": bb,
"ma": ma,
"equity": {
"dates": _eq_dates,
"strategy": _eq_strat,
"bh": _eq_bh,
},
"feature_importance": {
"names": [p[0] for p in fi_pairs],
"values": [float(p[1]) for p in fi_pairs],
},
"conf_matrix": cm.tolist(),
"conf_hist": {
"p_pos": _cp_pos,
"p_neg": _cp_neg,
},
"rolling_acc": {
"dates": _ra_dates,
"values": _ra_vals,
},
"drawdown": {
"dates": _dd_dates,
"values": _dd_vals,
},
"ret_dist": [float(x) for x in bt_result["trade_returns"]],
"ret_dist_long": [float(x) for x in bt_result["long_returns"]],
"ret_dist_short": [float(x) for x in bt_result["short_returns"]],
"train_ret_dist": [float(x) for x in bt_train_result["trade_returns"]] if bt_train_result else [],
"train_ret_dist_long": [float(x) for x in bt_train_result["long_returns"]] if bt_train_result else [],
"train_ret_dist_short": [float(x) for x in bt_train_result["short_returns"]] if bt_train_result else [],
"trade_log": _tl_test,
"train_trade_log": _tl_tr,
"trade_log_total": _tl_test_n,
"train_trade_log_total": _tl_tr_n,
"trade_log_truncated": _tl_test_tr,
"train_trade_log_truncated": _tl_tr_tr,
**(pre_stats or {}),
"metrics": metrics,
"split_dt": split_result["split_dt"],
"split_idx": int(split_result["split_idx"]),
"n_train": int(split_result["n_train"]),
"n_test": int(split_result["n_test"]),
"feature_cols": list(feature_cols),
"custom_figs": custom_figs or [],
}
# ════════════════════════════════════════════════════════════════════════════
# STRATEGY FRAMEWORK v2 — Config-driven architecture
# Claude writes feature_engineering() + strategy_config(). Framework does rest.
# ════════════════════════════════════════════════════════════════════════════
import importlib
_MODEL_REGISTRY = {
"XGBClassifier": ("xgboost", "XGBClassifier"),
"RandomForestClassifier": ("sklearn.ensemble", "RandomForestClassifier"),
"GradientBoostingClassifier": ("sklearn.ensemble", "GradientBoostingClassifier"),
"LogisticRegression": ("sklearn.linear_model", "LogisticRegression"),
"ExtraTreesClassifier": ("sklearn.ensemble", "ExtraTreesClassifier"),
"AdaBoostClassifier": ("sklearn.ensemble", "AdaBoostClassifier"),
}
RULE_MODEL_TYPE = "RuleBased"
def _build_model_from_config(config, X_train, y_train_enc):
"""Build, fit, and wrap a model from strategy_config dict."""
model_type = config.get("model_type", "RandomForestClassifier")
model_params = dict(config.get("model_params", {}))
# #374 — rule-based mode: a fixed human-written rule, NOT a fitted model.
# Nothing is trained; y_train_enc is deliberately ignored. predict_proba
# returns hard 0/1 so signal_threshold cannot alter a decision and reruns
# are bit-identical. Wrapped in ModelWrapper so the PKL/deploy/predictor
# path is byte-for-byte the same shape as a trained model's.
if model_type == RULE_MODEL_TYPE:
from model_wrapper import ModelWrapper, RuleModel
rules = config.get("rules") or []
if not rules:
raise ValueError(
"model_type 'RuleBased' requires a non-empty 'rules' list in "
"strategy_config() — e.g. "
"[{'signal': 1, 'when': [{'feature': 'rsi14', 'op': '<', 'value': 30}]}]"
)
rm = RuleModel(rules, n_features=X_train.shape[1])
rm.rule_signal(X_train.head(1)) # fail fast on a bad feature/op
return ModelWrapper(rm, original_classes=np.array([-1, 0, 1]),
n_features=X_train.shape[1])
if model_type not in _MODEL_REGISTRY:
raise ValueError(f"Unknown model_type '{model_type}'. "
f"Valid: {list(_MODEL_REGISTRY.keys()) + [RULE_MODEL_TYPE]}")
module_path, class_name = _MODEL_REGISTRY[model_type]
mod = importlib.import_module(module_path)
cls = getattr(mod, class_name)
# XGBoost defaults
if class_name == "XGBClassifier":
model_params.setdefault("use_label_encoder", False)
model_params.setdefault("eval_metric", "mlogloss")
model_params.setdefault("tree_method", "hist")
# Determinism > speed (2026-05-25). XGBoost hist with n_jobs=-1 is
# NON-reproducible even with random_state set — the parallel histogram
# gradient-sum order varies across threads, so the SAME code + data
# gives a slightly different model (and backtest) every run. Forcing
# single-thread makes training bit-reproducible so: (a) a user who
# copies a strategy and reruns it gets identical numbers, (b) the
# community "Live" score matches a redeploy, (c) "same code, different
# result" support reports go away. Cost: single-threaded XGB (a few
# seconds slower on large windows; hist is fast so it's minor). FORCED
# (not setdefault) so the guarantee can't be silently broken by a
# strategy passing n_jobs. Exact reproducibility holds within the
# platform (pinned versions / same Modal image); a user's own machine
# with different xgboost/numpy/CPU can still differ in low-order bits.
model_params["n_jobs"] = 1
# Common defaults
model_params.setdefault("random_state", 42)
from model_wrapper import ModelWrapper
clf = cls(**model_params)
# #405: label codes are deliberately FIXED at 0/1/2 for -1/0/1 so a code
# means the same thing in every model, run and deployment. The cost is
# that a window missing a class hands the fit a non-contiguous set - a
# real prod case, GBPUSD 1h with no flat bars gives {0, 2}. XGBoost
# rejects that outright ("Expected: [0 1], got [0 2]"); sklearn accepts it
# but then returns one proba column per PRESENT class, which downstream
# code - indexing by global class position - misreads without noticing.
# So: fit on a dense remap, and tell the wrapper what each column means.
_present = np.unique(np.asarray(y_train_enc))
if len(_present) < 2:
raise ValueError(
f"The training window contains only one label class "
f"({_present.tolist()}), so there is nothing for a classifier to "
f"separate. Widen the training window, or lower target_horizon / "
f"the labelling threshold so both up and down moves appear.")
_dense = {int(c): i for i, c in enumerate(_present)}
_y_fit = np.asarray([_dense[int(v)] for v in np.asarray(y_train_enc)])
clf.fit(X_train, _y_fit)
enc = LabelEncoder()
enc.fit([-1, 0, 1])
return ModelWrapper(clf, original_classes=enc.classes_,
n_features=X_train.shape[1],
class_positions=_present)
def _generate_signals(model, X, threshold):
"""Framework-owned signal generation. Deterministic threshold logic."""
proba = model.predict_proba(X)
classes = list(model.classes_)
idx_pos = classes.index(1) if 1 in classes else None
idx_neg = classes.index(-1) if -1 in classes else None
p_pos = proba[:, idx_pos] if idx_pos is not None else np.zeros(len(X))
p_neg = proba[:, idx_neg] if idx_neg is not None else np.zeros(len(X))
signal_vals = np.zeros(len(X))
signal_vals = np.where(p_pos >= threshold, 1.0, signal_vals)
signal_vals = np.where(p_neg >= threshold, -1.0, signal_vals)
# Both exceed: pick stronger
both = (p_pos >= threshold) & (p_neg >= threshold)
signal_vals[both] = np.where(p_pos[both] >= p_neg[both], 1.0, -1.0)
return pd.Series(signal_vals, index=X.index), p_pos, p_neg
# ── Filter functions (all no-ops when config value is None) ──────────────
def _apply_direction_filter(signal, direction):
"""Zero out signals that don't match allowed direction."""
if direction is None or direction == "both":
return signal
s = signal.copy()
if direction == "long":
s[s < 0] = 0.0
elif direction == "short":
s[s > 0] = 0.0
return s
def _apply_session_filter(signal, index, session_hours):
"""Zero out signals outside session hours [start, end] UTC."""
if session_hours is None:
return signal
s = signal.copy()
start_h, end_h = session_hours[0], session_hours[1]
hours = index.hour
if start_h <= end_h:
mask = (hours >= start_h) & (hours < end_h)
else: # wrap around midnight, e.g. [22, 6]
mask = (hours >= start_h) | (hours < end_h)
s[~mask] = 0.0
return s
def _apply_atr_filter(signal, close, high, low, min_atr):
"""Zero out signals when NATR(14) is below threshold."""
if min_atr is None:
return signal
hl = high - low
hc = (high - close.shift(1)).abs()
lc = (low - close.shift(1)).abs()
tr = pd.concat([hl, hc, lc], axis=1).max(axis=1)
atr14 = tr.ewm(com=13, adjust=False).mean()
natr = atr14 / close.replace(0, np.nan)
s = signal.copy()
s[natr < min_atr] = 0.0
return s
def _apply_trend_filter(signal, close, trend_filter):
"""Only allow signals aligned with trend. e.g. 'sma_50': longs above SMA, shorts below."""
if trend_filter is None:
return signal
# Parse: "sma_50" → SMA with period 50
parts = trend_filter.lower().replace("-", "_").split("_")
if len(parts) >= 2 and parts[0] in ("sma", "ema"):
period = int(parts[1])
else:
return signal # unknown filter, skip
if parts[0] == "sma":
trend_line = close.rolling(period).mean()
else:
trend_line = close.ewm(span=period, adjust=False).mean()
s = signal.copy()
# Longs only above trend, shorts only below
s[(s > 0) & (close < trend_line)] = 0.0
s[(s < 0) & (close > trend_line)] = 0.0
return s
# ── run_backtest_v2: framework-owned SL/TP/cooldown/position management ──
def run_backtest_v2(signal, close, high, low, config, capital=10000, cost=2e-5, leverage=1.0,
lot_size=None, balance=None, symbol="",
commission_usd=None, swap_long_usd=None, swap_short_usd=None,
spread_pips=None):
"""Backtest with SL/TP/cooldown/direction handling built into the engine.
Unlike run_backtest (v1), this function handles position exits internally.
Two accounting models (see run_backtest for the full contract):
• `lot_size is None` → LEGACY leverage×fractional return-on-equity; `pnl`
is the raw price delta. Direct callers (live_track / Server G) stay here.
• `lot_size` set → BROKER-ACCURATE: per-trade `pnl_usd` (position size
× per-pair FX − commission/swap); `pnl_pct` = pnl_usd / balance.
Returns: same dict shape as run_backtest()
"""
def _lev(r):
# Clamp at -100% so a single trade can't drive equity negative / flip
# the cumprod sign (SL'd trades are bounded; only un-SL'd adverse
# closes could approach the clamp).
return max(r * leverage, -1.0)
stop_loss = config.get("stop_loss")
take_profit = config.get("take_profit")
cooldown = config.get("cooldown", 0)
on_opposite = config.get("on_opposite", "reverse")
# FACT9 Phase 2 — SL/TP unit. 'pct' (legacy/default) = fraction of price
# (entry × level). 'pips' = pip distance (level × pip_size). 'usd' = the
# price move whose pre-cost USD P&L equals `level` dollars at this lot. The
# pips/usd modes need a position size → only meaningful in broker mode
# (lot_size set); 'pct' reproduces the original entry × (1 ± level) exactly.
risk_unit = (config.get("risk_unit") or "pct").lower()
_pip_sz = _pip_size(symbol)
_risk_units = (lot_size or 0.0) * LOT_UNITS
# Bid/ask spread cost (#283). Brokers fill buys at ask / sells at bid, so a
# round-trip pays one full spread; the engine otherwise fills at the raw
# close (mid), silently overstating P&L. `spread_pips` = full round-trip
# spread in pips, charged once per closed trade. Default 0.0 → no change to
# existing backtests. Realistic EURUSD: ~0.1 pip raw/ECN, ~1.0 pip retail.
# Precedence: explicit spread_pips arg > config["spread_pips"] > 0.0.
_sp = spread_pips if spread_pips is not None else config.get("spread_pips")
_spread_px = float(_sp or 0.0) * _pip_sz
def _risk_dist(level, e_px):
"""SL/TP price distance for `level` at entry `e_px`, per risk_unit."""
if level is None:
return None
if risk_unit == "pips":
return level * _pip_sz
if risk_unit == "usd":
return (level * _usd_conv(symbol, e_px) / _risk_units) if _risk_units > 0 else e_px * level
return e_px * level # 'pct' (legacy)
sig_arr = signal.values
close_arr = close.values
high_arr = high.values
low_arr = low.values
idx = signal.index
n = len(close_arr)
broker = lot_size is not None
bal = float(balance) if balance is not None else float(capital)
if broker and bal <= 0:
bal = float(capital) or 10000.0
_idx_dates, _idx_dow = _index_date_arrays(idx) if broker else (None, None)
_base = bal if broker else float(capital)
trade_returns = []
long_returns = []
short_returns = []
trade_log = []
equity_vals = np.full(n, _base)
cumret = 1.0
position = 0.0 # current direction: 1.0, -1.0, or 0.0 (flat)
entry_price = None
entry_bar = None # index into arrays for entry time
cooldown_remaining = 0
def _close_ret(direction, e_px, x_px, e_bar, x_bar):
"""Return (ret_for_equity, pnl_usd_or_None) for a closed trade."""
if broker:
nights = _nights_held(_idx_dates, _idx_dow, e_bar, x_bar)
usd = _trade_pnl_usd(direction, e_px, x_px, lot_size, symbol, nights,
commission_usd=commission_usd,
swap_long_usd=swap_long_usd,
swap_short_usd=swap_short_usd)
# one full round-trip spread (buy@ask + sell@bid), in USD (#283)
usd -= _spread_px * _risk_units / _usd_conv(symbol, x_px)
return max(usd / bal, -1.0), usd
# legacy leverage mode: subtract spread as a fraction of entry price (#283)
return _lev(float(direction * (x_px - e_px) / e_px - cost - _spread_px / e_px)), None
def _log_trade(exit_bar, exit_px, ret, reason, usd=None):
trade_log.append({
"type": "Buy" if position == 1.0 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[exit_bar]),
"entry_price": round(entry_price, 5),
"exit_price": round(exit_px, 5),
"pnl": round(position * (exit_px - entry_price), 5),
"pnl_usd": round(usd, 2) if usd is not None else None,
"pips": round(_pips(symbol, position, entry_price, exit_px), 1) if broker else None,
"pnl_pct": round(ret * 100, 3),
"exit_reason": reason,
})
for i in range(n):
c = close_arr[i]
h = high_arr[i]
lo = low_arr[i]
s = sig_arr[i]
# 1. Check SL/TP if in trade
if position != 0.0 and entry_price is not None:
hit_sl = False
hit_tp = False
exit_price = None
d_sl = _risk_dist(stop_loss, entry_price)
d_tp = _risk_dist(take_profit, entry_price)
if position == 1.0: # long
if d_sl is not None and lo <= entry_price - d_sl:
hit_sl = True
exit_price = entry_price - d_sl
elif d_tp is not None and h >= entry_price + d_tp:
hit_tp = True
exit_price = entry_price + d_tp
else: # short
if d_sl is not None and h >= entry_price + d_sl:
hit_sl = True
exit_price = entry_price + d_sl
elif d_tp is not None and lo <= entry_price - d_tp:
hit_tp = True
exit_price = entry_price - d_tp
if hit_sl or hit_tp:
ret, _usd = _close_ret(position, entry_price, exit_price, entry_bar, i)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, exit_price, ret, "SL" if hit_sl else "TP", usd=_usd)
cumret *= (1 + ret)
position = 0.0
entry_price = None
entry_bar = None
cooldown_remaining = cooldown
equity_vals[i] = _base * cumret
continue
# 2. Cooldown
if cooldown_remaining > 0:
cooldown_remaining -= 1
equity_vals[i] = _base * cumret
continue
# 3. Signal processing
if s != 0.0:
if position == 0.0:
# Open new trade
position = s
entry_price = c
entry_bar = i
elif s != position:
# Opposite signal
if on_opposite == "reverse":
# Close current + open opposite
ret, _usd = _close_ret(position, entry_price, c, entry_bar, i)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, c, ret, "signal", usd=_usd)
cumret *= (1 + ret)
position = s
entry_price = c
entry_bar = i
else: # close_only
# Close current, go flat
ret, _usd = _close_ret(position, entry_price, c, entry_bar, i)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, c, ret, "close_only", usd=_usd)
cumret *= (1 + ret)
position = 0.0
entry_price = None
entry_bar = None
cooldown_remaining = cooldown
equity_vals[i] = _base * cumret
# Close last open trade at final close
if position != 0.0 and entry_price is not None and n > 0 and entry_price != 0:
c = close_arr[-1]
ret, _usd = _close_ret(position, entry_price, c, entry_bar, n - 1)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(n - 1, c, ret, "end", usd=_usd)
cumret *= (1 + ret)
equity_vals[-1] = _base * cumret
# Bar returns for Sharpe (approximate)
bar_returns = np.zeros(n)
for i in range(1, n):
if close_arr[i - 1] != 0 and sig_arr[i - 1] != 0:
bar_returns[i] = sig_arr[i - 1] * (close_arr[i] - close_arr[i - 1]) / close_arr[i - 1]
return {
"equity": pd.Series(equity_vals, index=close.index),
"trade_returns": trade_returns,
"long_returns": long_returns,
"short_returns": short_returns,
"bar_returns": bar_returns,
"trade_log": trade_log,
}
# ── run_strategy: the v2 orchestrator ────────────────────────────────────
def run_strategy(feature_fn, config_fn, data_path, start_date="", end_date="",
validation_date="", train_split=0.7, register_model_fn=None,
leverage=30.0, lot_size=None, balance=None, symbol="",
stop_loss=None, take_profit=None, risk_unit=None,
signal_threshold=None,
commission_usd=None, swap_long_usd=None, swap_short_usd=None,
spread_pips=None, aux_feeds=None,
purge_bars=0, embargo_bars=0):
"""Config-driven strategy execution. Claude writes feature_fn + config_fn,
framework does everything else.
`lot_size`/`balance`/`symbol` drive the broker-accurate $ accounting (Phase 1
Factual Trading Product). `lot_size is None` keeps the LEGACY leverage path —
direct callers that omit it (e.g. live_track / Server G) are unchanged. The
synthesized train wrapper passes lot_size=1.0 by default → broker mode.
`aux_feeds` (#387/#388): explicit list of {"symbol","timeframe"} dicts, or
None to fall back to config_fn()["aux_feeds"]. An explicit [] means "no
aux feeds" and does NOT fall through to config. Normalized against the
RUN-TIME primary (symbol/data_path), which is how a pasted strategy whose
aux feed collides with a dropdown-selected primary gets that feed dropped.
`purge_bars`/`embargo_bars` (#411): passed straight through to
split_bounds() at the train/test split below. Defaults (0, 0) reproduce
the legacy split exactly - the headline 70/30 backtest is unaffected.
Returns: results dict (same format as webapp expects)
"""
config = config_fn()
# #244 — an explicitly-passed confidence/signal threshold (API
# confidence_threshold or Config tab) is the source of truth: override the
# generated strategy_config() value. None = strategy's own value stands.
if signal_threshold is not None:
config["signal_threshold"] = float(signal_threshold)
# FACT9 Phase 2 — Config-tab SL/TP are the source of truth: override the
# generated strategy_config() when the caller (synthesized wrapper) provides
# them. risk_unit selects interpretation downstream (pct/pips/usd). In
# pips/usd mode Config OWNS the levels — any pct-scale value strategy_config()
# left behind is dropped (a 0.005 ratio must never be read as 0.005 pips).
if risk_unit is not None:
_ru = str(risk_unit).lower()
config["risk_unit"] = _ru
if _ru in ("pips", "usd"):
config["stop_loss"] = stop_loss
config["take_profit"] = take_profit
else:
if stop_loss is not None:
config["stop_loss"] = stop_loss
if take_profit is not None:
config["take_profit"] = take_profit
else:
if stop_loss is not None:
config["stop_loss"] = stop_loss
if take_profit is not None:
config["take_profit"] = take_profit
# Broker-accounting params. symbol falls back to the SYMBOL parsed from the
# data path (e.g. .../USDJPY_15min.parquet) so per-pair FX conversion works
# even when the caller doesn't pass it explicitly.
_broker = lot_size is not None
_sym = symbol or config.get("symbol") or ""
if not _sym:
_sym, _ = _parse_symbol_tf_from_path(data_path)
_sym = _sym or ""
_bal = float(balance) if balance is not None else 10000.0
_cap = _bal if _broker else 10000
# Auto-correct SL/TP ONLY in legacy price-% mode (Claude passing 50 for 0.5).
# pips/usd units are absolute magnitudes (25 pips, $250) → never divide them.
if (config.get("risk_unit") or "pct").lower() == "pct":
for _key in ("stop_loss", "take_profit"):
_val = config.get(_key)
if _val is not None and _val > 0.1: # >10% is almost certainly a percentage
config[_key] = _val / 100.0
print(f"[strategy] Auto-corrected {_key}: {_val} -> {config[_key]} (was percentage, converted to decimal)")
# 1. Load data
df, close, open_, high, low = load_ohlc(data_path, start_date, end_date)
# 1b. #387/#388 - merge aux feeds BEFORE feature_engineering. Run-time
# normalization (against _sym/_tf, the RUN-TIME primary) is where a
# pasted strategy whose aux feed now equals the dropdown-selected
# primary gets that feed dropped - one rule, no special case at paste
# time. Precedence: explicit aux_feeds kwarg (incl. []) wins; None falls
# back to config_fn()["aux_feeds"].
_aux_raw = aux_feeds if aux_feeds is not None else config.get("aux_feeds")
_tf = _parse_symbol_tf_from_path(data_path)[1] or "15min"
_aux_specs, _aux_warns = normalize_aux_specs(_aux_raw or [], _sym, _tf)
df, _aux_report = merge_aux_feeds(df, _aux_specs, _sym, _tf)
for _w in _aux_warns:
_aux_report["dropped"].append({"reason": _w})
# #403: say it out loud. A dropped feed is caught into the report and the
# run continues (correct - a missing aux feed must not kill a train), but
# with nothing printed, an aux feature that never loaded on Modal looked
# exactly like a successful run for a full day. The log line is the only
# thing standing between "degraded" and "undetectably broken".
if _aux_report["dropped"]:
for _d in _aux_report["dropped"]:
print(f"[strategy] aux feed DROPPED: {_d}")
if _aux_report["merged"]:
print(f"[strategy] aux feeds merged: {_aux_report['merged']}")
# merge_aux_feeds returns a (possibly new, e.g. via sort_index) df object
# when aux feeds were merged - refresh close/open_/high/low so they read
# from the SAME frame handed to feature_fn rather than stale references
# into the pre-merge object (values are unaffected either way - the
# merge never touches primary OHLC or drops rows - but identity must
# match, not just values, for feature_fn to see a self-consistent df).
close = df["close"]
open_ = df["open"]
high = df["high"]
low = df["low"]
# 2. Feature engineering (Claude's function)
df = feature_fn(df, close, open_, high, low)
close = df["close"]
open_ = df["open"]
high = df["high"]
low = df["low"]
# 3. Warm-up detection: drop rows where features have NaN BEFORE any fill
feature_cols = [c for c in df.columns if c not in ("open", "high", "low", "close")]
raw_nans = df[feature_cols].isna().any(axis=1)
valid_rows = ~raw_nans
if valid_rows.any():
first_valid = valid_rows.idxmax()
if raw_nans.loc[:first_valid].any():
df = df.loc[first_valid:].copy()
close = df["close"]
open_ = df["open"]
high = df["high"]
low = df["low"]
# 4. Target
horizon = config.get("target_horizon", 4)
target = make_target(close, horizon=horizon)
# 5. Split (ffill only within each partition — no bfill leak)
mask = target.notna()
df = df[mask].copy()
target = target[mask]
close = df["close"]
high = df["high"]
low = df["low"]
X = df[feature_cols].copy()
X = X.replace([np.inf, -np.inf], np.nan)
if validation_date:
split_idx = len(df[df.index <= validation_date])
else:
split_idx = int(len(df) * train_split)
# #411 - the ONE chronological splitter. purge_bars trims the tail of
# train so no row's forward-looking label (make_target) reaches into
# test; embargo_bars pushes test_start out further. Defaults (0, 0)
# collapse train_end == test_start == split_idx, i.e. bit-identical to
# the legacy split.
train_end, test_start = split_bounds(len(df), split_idx,
purge_bars, embargo_bars)
split_idx = test_start # split_dt / n_test keep their existing meaning
# ffill within train and test separately (no leak)
X_train = X.iloc[:train_end].ffill().fillna(0.0)
X_test = X.iloc[test_start:].ffill().fillna(0.0)
X = pd.concat([X_train, X_test])
y_train = target.iloc[:train_end]
y_test = target.iloc[test_start:]
close_train = close.iloc[:train_end]
close_test = close.iloc[test_start:]
high_test = high.iloc[test_start:]
low_test = low.iloc[test_start:]
enc = LabelEncoder()
enc.fit([-1, 0, 1])
y_train_enc = enc.transform(y_train)
y_test_enc = enc.transform(y_test)
split_dt = str(df.index[split_idx])
sp = {
"df": df, "X_train": X_train, "X_test": X_test,
"y_train": y_train, "y_test": y_test,
"y_train_enc": y_train_enc, "y_test_enc": y_test_enc,
"enc": enc,
"close": close, "close_train": close_train, "close_test": close_test,
"split_idx": split_idx, "split_dt": split_dt,
"n_train": len(X_train), "n_test": len(X_test),
}
# 6. Build model from config
model = _build_model_from_config(config, X_train, y_train_enc)
# 7. Generate signals
threshold = config.get("signal_threshold", 0.55)
signal_train, p_pos_train, p_neg_train = _generate_signals(model, X_train, threshold)
signal_test, p_pos_test, p_neg_test = _generate_signals(model, X_test, threshold)
# 8. Apply filters (order: direction → session → ATR → trend)
direction = config.get("direction", "both")
signal_test = _apply_direction_filter(signal_test, direction)
signal_train = _apply_direction_filter(signal_train, direction)
session_filter = config.get("session_filter")
signal_test = _apply_session_filter(signal_test, signal_test.index, session_filter)
signal_train = _apply_session_filter(signal_train, signal_train.index, session_filter)
min_atr = config.get("min_atr")
if min_atr is not None:
signal_test = _apply_atr_filter(signal_test, close_test, high_test, low_test, min_atr)
trend_filter = config.get("trend_filter")
if trend_filter is not None:
signal_test = _apply_trend_filter(signal_test, close_test, trend_filter)
signal_full = pd.concat([signal_train, signal_test])
# 9. Backtest with SL/TP/cooldown (test + train)
high_train = high.iloc[:train_end]
low_train = low.iloc[:train_end]
has_risk = (config.get("stop_loss") is not None or
config.get("take_profit") is not None or
config.get("cooldown", 0) > 0 or
config.get("on_opposite", "reverse") != "reverse")
# Bid/ask spread (#283): caller value wins, else a per-symbol default so
# every training applies an honest spread (FX majors ~0.15 pip; crypto 0.0).
_spread = float(spread_pips) if spread_pips is not None else _default_spread_pips(_sym)
_costs = dict(commission_usd=commission_usd, swap_long_usd=swap_long_usd,
swap_short_usd=swap_short_usd, spread_pips=_spread)
if has_risk:
bt = run_backtest_v2(signal_test, close_test, high_test, low_test, config, capital=_cap, leverage=leverage, lot_size=lot_size, balance=_bal, symbol=_sym, **_costs)
bt_train = run_backtest_v2(signal_train, close_train, high_train, low_train, config, capital=_cap, leverage=leverage, lot_size=lot_size, balance=_bal, symbol=_sym, **_costs)
else:
bt = run_backtest(signal_test, close_test, capital=_cap, leverage=leverage, lot_size=lot_size, balance=_bal, symbol=_sym, **_costs)
bt_train = run_backtest(signal_train, close_train, capital=_cap, leverage=leverage, lot_size=lot_size, balance=_bal, symbol=_sym, **_costs)
# 10. Metrics
metrics = compute_metrics(bt, close_test, capital=_cap)
# 11. Pre-compute all trade stats (single source of truth)
pre_stats = {
"train_stats": compute_trade_stats(bt_train.get("trade_returns", []), capital=_cap),
"test_stats": compute_trade_stats(bt.get("trade_returns", []), capital=_cap),
"long_stats": compute_trade_stats(bt.get("long_returns", []), capital=_cap),
"short_stats": compute_trade_stats(bt.get("short_returns", []), capital=_cap),
}
# 12. Register model
if register_model_fn is not None:
register_model_fn(model)
# 13. Build return dict
_results = build_return_dict(sp, bt, metrics, model, feature_cols,
signal_full, p_pos_test, p_neg_test, custom_figs=[],
bt_train_result=bt_train, pre_stats=pre_stats)
# #387/#388 - aux feed merge/drop report (incl. normalize warnings).
_results["aux_feeds_report"] = _aux_report
# Stash the leverage actually used so display/predictor/API can read it.
_results["leverage"] = leverage
# Broker-accounting params (Phase 1) — None lot_size = legacy path.
if _broker:
_results["lot_size"] = lot_size
_results["balance"] = _bal
_results["symbol"] = _sym
# FACT9 Phase 2 — risk model (unit + levels) + a FIXED pip-equivalent for the
# live predictor (Server B is HTTP-only → can't convert $ live; the backtest
# already applied the exact per-trade $, live uses this constant pip stop).
_ru2 = (config.get("risk_unit") or "pct").lower()
_results["risk_unit"] = _ru2
_results["stop_loss"] = config.get("stop_loss")
_results["take_profit"] = config.get("take_profit")
if _ru2 in ("pips", "usd"):
try:
_refpx = float(close.iloc[-1])
_ps2 = _pip_size(_sym)
_u2 = (lot_size or 1.0) * LOT_UNITS
def _lvl_to_pips(_lv):
if _lv is None:
return None
if _ru2 == "pips":
return float(_lv)
return float(_lv) * _usd_conv(_sym, _refpx) / (_ps2 * _u2) # usd→pips
_results["sl_pips"] = _lvl_to_pips(config.get("stop_loss"))
_results["tp_pips"] = _lvl_to_pips(config.get("take_profit"))
except Exception:
pass
return _results
# ── End strategy_utils ──
DATA_PATH = '/root/Desktop/QuantifyMe/data/ohlc/EURUSD_15min.parquet'
START_DATE = '2025-07-31'
END_DATE = '2026-02-03'
VALIDATION_DATE = ""
TRAIN_SPLIT = 0.6999526739233317
LEVERAGE = 30.0
LOTS = 1.0
BALANCE = 10000.0
RISK_UNIT = 'pips'
STOP_LOSS = 25.0
TAKE_PROFIT = 50.0
def feature_engineering(df, close, open_, high, low):
"""
Feature Engineering: Standard momentum & volatility indicators
- Multi-horizon returns: 1, 3, 5, 10, 20 bars
- RSI 14: momentum oscillator
- Bollinger Bands 20-period, 2 std: volatility & mean reversion
- Moving averages 50, 100, 200: trend structure
No lookahead bias. NaN filled via bfill/ffill at end.
"""
# ============ Multi-horizon Returns ============
df['ret_1'] = close.pct_change(1)
df['ret_3'] = close.pct_change(3)
df['ret_5'] = close.pct_change(5)
df['ret_10'] = close.pct_change(10)
df['ret_20'] = close.pct_change(20)
# ============ RSI 14 ============
delta = close.diff()
gain = np.where(delta > 0, delta, 0)
loss = np.where(delta < 0, -delta, 0)
gain_ema = pd.Series(gain, index=close.index).ewm(span=14, adjust=False).mean()
loss_ema = pd.Series(loss, index=close.index).ewm(span=14, adjust=False).mean()
rs = gain_ema / (loss_ema + 1e-10)
df['rsi_14'] = 100 - (100 / (1 + rs))
# ============ Bollinger Bands 20, 2 ============
sma_20 = close.rolling(20).mean()
std_20 = close.rolling(20).std()
df['bb_upper'] = sma_20 + 2 * std_20
df['bb_lower'] = sma_20 - 2 * std_20
df['bb_mid'] = sma_20
df['bb_width'] = df['bb_upper'] - df['bb_lower']
df['bb_pct'] = (close - df['bb_lower']) / (df['bb_width'] + 1e-10)
# ============ Moving Averages (Trend Filter) ============
df['sma_50'] = close.rolling(50).mean()
df['sma_100'] = close.rolling(100).mean()
df['sma_200'] = close.rolling(200).mean()
# ============ ATR for Volatility (used in framework filters) ============
tr = np.maximum.reduce([
high - low,
np.abs(high - close.shift(1)),
np.abs(low - close.shift(1))
])
df['atr_14'] = pd.Series(tr, index=close.index).rolling(14).mean()
df['natr_14'] = df['atr_14'] / close * 100
# ============ Fill NaN from indicator warm-up ============
df = df.bfill().ffill()
return df
def strategy_config():
"""
XGBoost classifier optimized for Sharpe ratio on EUR/USD 15-min data.
Strategy:
- Model: XGBClassifier (binary classification: up/down next 4 bars)
- Signal: confidence (predict_proba) > 0.55
- Direction: both (long and short)
- Position: 1 concurrent trade, close_only on opposite signal
- Risk: no hard SL/TP at config level; framework applies cooldown & position mgmt
Hyperparameters tuned for balance between:
* Sharpe ratio (high-quality signals, low noise)
* Drawdown control (regularization + shallow trees)
* Transaction cost tolerance (at 2e-5 round-trip)
"""
return {
"title": "EUR/USD XGBoost RSI + Bollinger Bands Scalper",
"model_type": "XGBClassifier",
"model_params": {
"n_estimators": 300,
"max_depth": 5,
"learning_rate": 0.05,
"subsample": 0.8,
"colsample_bytree": 0.8,
"gamma": 1.0,
"min_child_weight": 2,
"random_state": 42,
},
"signal_threshold": 0.55,
"direction": "both",
"stop_loss": None,
"take_profit": None,
"cooldown": 0,
"max_positions": 1,
"on_opposite": "close_only",
"session_filter": None,
"min_atr": None,
"trend_filter": None,
"target_horizon": 4,
"objective": (
"Maximize Sharpe ratio on EUR/USD 15-min data via XGBClassifier. "
"Features: returns (1, 3, 5, 10, 20 bars), RSI 14, Bollinger Bands 20/2, "
"SMAs (50, 100, 200), ATR 14. Model trained on 70% historical data, "
"tested on 30% holdout (2025-07-31 to 2026-07-31)."
),
"notes": (
"Regularization (gamma=1.0, min_child_weight=2, max_depth=5) keeps tree complexity low "
"to improve generalization and reduce overfitting. Learning rate 0.05 with 300 trees "
"provides stable gradient descent. Subsample & colsample at 0.8 add robustness. "
"Signal threshold 0.55 filters weak predictions. No hard SL/TP; framework handles "
"risk via position management. Target horizon 4 bars (~1 hour on 15-min) balances "
"scalper speed with signal clarity."
),
}
# ---------------------------------------------------------------
# AUX_FEEDS - written by the strategy BUILDER, not by the model.
# Last module-level statement on purpose: module execution is
# last-write-wins, so this is the value run_strategy() sees no matter
# what appears above. Do not move it, and add nothing after it.
# ---------------------------------------------------------------
AUX_FEEDS = []
# ── Framework v2: auto-generated wrapper ──
def train_and_backtest():
_vd = VALIDATION_DATE if 'VALIDATION_DATE' in globals() else ''
_ts = TRAIN_SPLIT if 'TRAIN_SPLIT' in globals() else 0.7
_aux = AUX_FEEDS if 'AUX_FEEDS' in globals() else None
_lev = LEVERAGE if 'LEVERAGE' in globals() else 30.0
_lots = LOTS if 'LOTS' in globals() else 1.0
_bal = BALANCE if 'BALANCE' in globals() else 10000.0
_ru = RISK_UNIT if 'RISK_UNIT' in globals() else None
_sl = STOP_LOSS if 'STOP_LOSS' in globals() else None
_tp = TAKE_PROFIT if 'TAKE_PROFIT' in globals() else None
_sth = SIGNAL_THRESHOLD if 'SIGNAL_THRESHOLD' in globals() else None
_comm = QM_COMMISSION_USD if 'QM_COMMISSION_USD' in globals() else None
_swl = QM_SWAP_LONG_USD if 'QM_SWAP_LONG_USD' in globals() else None
_sws = QM_SWAP_SHORT_USD if 'QM_SWAP_SHORT_USD' in globals() else None
return run_strategy(
feature_engineering, strategy_config,
DATA_PATH, START_DATE, END_DATE,
_vd, _ts,
register_model_fn=register_model,
leverage=_lev, lot_size=_lots, balance=_bal,
stop_loss=_sl, take_profit=_tp, risk_unit=_ru,
signal_threshold=_sth,
commission_usd=_comm, swap_long_usd=_swl, swap_short_usd=_sws,
aux_feeds=_aux
)
|
||||||||||
|
🥉
|
EMA(9/21) trend
|
M
@malcolmtan
|
EURUSD | 1min | 50.0%33.3% | +0.90%-1.58% | 2.150.71 | 0.39%0.39% | 1015 |
|
# ╔══════════════════════════════════════════════════════════════╗
# ║ STRATEGY REQUEST LOG ║
# ╚══════════════════════════════════════════════════════════════╝
# Generated : 2026-05-25 02:27:36
# Model : XGBoost
# Feature Eng. : go long when EMA(9) crosses above EMA(21), exit when it crosses back below + Auto-add features: ON
# Signal / Entry : —
# Optimization : —
# Risk Mgmt : —
# Risk Filter : —
# ══════════════════════════════════════════════════════════════
# ============================================================
# SECTION 0 — IMPORTS & CONSTANTS
import numpy as np
import pandas as pd
# ── Inlined strategy_utils ──
"""
strategy_utils.py — Standard utility functions for generated strategies.
Claude imports these instead of writing boilerplate from scratch.
This ensures consistent behavior across all generated strategies.
"""
import numpy as np
import pandas as pd
from sklearn.preprocessing import LabelEncoder
# Max backtest window per timeframe. A finer timeframe over a longer window
# blows up the results dict / parquet load / Modal train time (the 2026-05-12
# OOM was a 1-min × multi-year sweep) — and a 1-min strategy gains nothing from
# 2 years of 1-min bars. Enforced HERE because every training path (UI / API /
# Modal) funnels through run_strategy → load_ohlc. Env-overridable so a future
# "max plan" / dedicated-server tier can lift it.
_TF_MAX_DAYS = {
"1min": 30,
"5min": 90,
"15min": 365,
"1h": 730,
}
def _fetch_ohlc_from_internal(symbol: str, tf: str, start: str, end: str):
"""Phase 3.2: fetch parquet bytes from Server A's /internal/ohlc endpoint
instead of reading a local file. Used inside Modal containers / Mac worker
pool (Phase 3.4) so every train sees the same source of truth as the chart.
Returns: pd.DataFrame (parquet decoded), or raises on any failure so the
caller can fall back / surface a clear error in the job.
"""
import hashlib as _hashlib, hmac as _hmac, io as _io, os as _os
import urllib.request as _ur, urllib.parse as _urp
base = (_os.environ.get("QM_INTERNAL_OHLC_BASE") or "").rstrip("/")
secret = (_os.environ.get("INTERNAL_WS_SECRET") or "").strip()
if not base:
raise RuntimeError("QM_INTERNAL_OHLC_BASE not set")
if not secret:
raise RuntimeError("INTERNAL_WS_SECRET not set")
msg = f"{symbol}|{tf}|{start}|{end}".encode("utf-8")
sig = _hmac.new(secret.encode("utf-8"), msg, _hashlib.sha256).hexdigest()
qs = _urp.urlencode({
"symbol": symbol, "tf": tf,
"start": start, "end": end, "sig": sig,
})
url = f"{base}/internal/ohlc?{qs}"
req = _ur.Request(url, headers={"User-Agent": "qm-worker/1.0"})
with _ur.urlopen(req, timeout=30) as resp:
if resp.status != 200:
raise RuntimeError(f"/internal/ohlc returned {resp.status}")
payload = resp.read()
print(f"[load_ohlc:internal] {symbol} {tf} fetched {len(payload)} bytes", flush=True)
return pd.read_parquet(_io.BytesIO(payload))
def _parse_symbol_tf_from_path(data_path: str):
"""Pull SYMBOL + TF out of a path like .../EURUSD_1min.parquet."""
import os as _os, re as _re
base = _os.path.basename(str(data_path))
m = _re.match(r"^([A-Z]{6})_(\d+min|\d+h)\.parquet$", base)
if not m:
return None, None
return m.group(1), m.group(2)
def load_ohlc(data_path, start_date="", end_date=""):
"""Load OHLC parquet, sort index, filter dates. Always returns consistent format.
The lower bound is clamped per timeframe (see _TF_MAX_DAYS) — a request for
more history than the cap silently starts later.
Phase 3.2: when env QM_USE_INTERNAL_OHLC=="1", fetch over HTTP from
Server A's /internal/ohlc endpoint instead of pd.read_parquet on a local
file (which on Modal is a stale Volume snapshot). The endpoint applies the
same day-cap, so the local cap-check below is a defensive no-op in that
path. Flag defaults to "0" → unchanged behavior.
Returns: (df, close, open_, high, low)
"""
import os as _os, re as _re
_use_internal = _os.environ.get("QM_USE_INTERNAL_OHLC", "0") == "1"
if _use_internal:
_sym, _tf = _parse_symbol_tf_from_path(data_path)
if not _sym or not _tf:
raise RuntimeError(
f"QM_USE_INTERNAL_OHLC=1 but DATA_PATH basename does not match "
f"SYMBOL_TF.parquet: {data_path}"
)
df = _fetch_ohlc_from_internal(_sym, _tf, start_date or "", end_date or "")
else:
df = pd.read_parquet(data_path)
df.index = pd.to_datetime(df.index)
df = df.sort_index()
# Per-timeframe window cap (timeframe inferred from the parquet filename).
_m = _re.search(r"_(\d+min|\d+h)\.parquet$", _os.path.basename(str(data_path)))
_tf = _m.group(1) if _m else None
_max_days = _TF_MAX_DAYS.get(_tf)
if _max_days and _max_days > 0 and len(df):
_env_override = _os.environ.get(f"QM_MAX_DAYS_{_tf.upper()}")
if _env_override and _env_override.isdigit():
_max_days = int(_env_override)
try:
_eff_end = pd.Timestamp(end_date) if end_date else df.index.max()
_eff_end = min(_eff_end, df.index.max())
_floor = _eff_end - pd.Timedelta(days=_max_days)
_req_start = pd.Timestamp(start_date) if start_date else df.index.min()
if _req_start < _floor:
print(f"[load_ohlc] {_tf} backtest window capped to {_max_days}d: "
f"start {_req_start.date()} -> {_floor.date()}", flush=True)
start_date = _floor
except Exception as _e:
print(f"[load_ohlc] window-cap check skipped ({_e})", flush=True)
if start_date:
df = df[df.index >= start_date]
if end_date:
df = df[df.index <= end_date]
return df, df["close"], df["open"], df["high"], df["low"]
def make_target(close, horizon=4):
"""Create target: direction N bars ahead. Default 4 bars = 1 hour on 15-min data.
Returns: target (pd.Series of -1, 0, 1)
"""
return np.sign(close.shift(-horizon) - close)
def split_data(df, target, feature_cols, train_split=0.7, validation_date=""):
"""Train/test split. Handles both ratio and date-based splits.
Drops NaN from target before splitting. Encodes labels to [0,1,2].
Returns: dict with keys:
X_train, X_test, y_train, y_test,
y_train_enc, y_test_enc, enc,
close_train, close_test,
split_idx, split_dt, n_train, n_test
"""
# Drop NaN from target
mask = target.notna()
df = df[mask].copy()
target = target[mask]
close = df["close"]
# Build feature matrix
X = df[feature_cols].copy()
X = X.bfill().ffill()
X = X.replace([np.inf, -np.inf], np.nan).fillna(0.0)
# Split
if validation_date:
split_idx = len(df[df.index <= validation_date])
else:
split_idx = int(len(df) * train_split)
split_idx = max(1, min(split_idx, len(df) - 1))
X_train = X.iloc[:split_idx]
X_test = X.iloc[split_idx:]
y_train = target.iloc[:split_idx]
y_test = target.iloc[split_idx:]
close_train = close.iloc[:split_idx]
close_test = close.iloc[split_idx:]
split_dt = str(df.index[split_idx])
# Label encoding — always fit on [-1, 0, 1]
enc = LabelEncoder()
enc.fit([-1, 0, 1])
y_train_enc = enc.transform(y_train)
y_test_enc = enc.transform(y_test)
return {
"df": df, "X_train": X_train, "X_test": X_test,
"y_train": y_train, "y_test": y_test,
"y_train_enc": y_train_enc, "y_test_enc": y_test_enc,
"enc": enc,
"close": close, "close_train": close_train, "close_test": close_test,
"split_idx": split_idx, "split_dt": split_dt,
"n_train": len(X_train), "n_test": len(X_test),
}
def compute_overlays(close, df_index):
"""Compute BB and MA overlays on full dataset. Always consistent.
Returns: (bb_dict, ma_dict)
"""
bb_mid = close.rolling(20).mean()
bb_std = close.rolling(20).std()
bb_upper = bb_mid + 2 * bb_std
bb_lower = bb_mid - 2 * bb_std
ma50 = close.rolling(50).mean()
ma100 = close.rolling(100).mean()
ma200 = close.rolling(200).mean()
def _safe(s):
s = s.reindex(df_index).bfill().ffill()
return [float(x) if (x is not None and not np.isnan(x) and not np.isinf(x)) else None
for x in s.values]
bb = {"upper": _safe(bb_upper), "mid": _safe(bb_mid), "lower": _safe(bb_lower)}
ma = {"ma50": _safe(ma50), "ma100": _safe(ma100), "ma200": _safe(ma200)}
return bb, ma
def run_backtest(signal, close, capital=10000, cost=2e-5):
"""Run backtest with transaction costs.
Uses price-based trade returns (same as webapp _compute_trades).
Signal 0 = hold (keep current position), not close.
Returns: dict with equity, trade_returns, long_returns, short_returns, bar_returns
"""
sig_arr = signal.values
price_arr = close.values
idx = signal.index
n = len(price_arr)
# Trade returns — price-based (matches webapp _compute_trades exactly)
trade_returns = []
long_returns = []
short_returns = []
trade_log = []
last_dir = None
entry_price = None
entry_bar = None
for i in range(n):
s = sig_arr[i]
c = price_arr[i]
if s != 0.0 and s != last_dir:
# Direction change — close previous trade, open new
if last_dir is not None and entry_price is not None and entry_price != 0:
ret = float(last_dir * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if last_dir == 1:
long_returns.append(ret)
else:
short_returns.append(ret)
trade_log.append({
"type": "Buy" if last_dir == 1 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[i]),
"entry_price": round(entry_price, 5),
"exit_price": round(c, 5),
"pnl": round(last_dir * (c - entry_price), 5),
"pnl_pct": round(ret * 100, 3),
"exit_reason": "signal",
})
entry_price = c
entry_bar = i
last_dir = s
# Close last open trade
if last_dir is not None and entry_price is not None and n > 0 and entry_price != 0:
c = price_arr[-1]
ret = float(last_dir * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if last_dir == 1:
long_returns.append(ret)
else:
short_returns.append(ret)
trade_log.append({
"type": "Buy" if last_dir == 1 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[-1]),
"entry_price": round(entry_price, 5),
"exit_price": round(c, 5),
"pnl": round(last_dir * (c - entry_price), 5),
"pnl_pct": round(ret * 100, 3),
"exit_reason": "end",
})
# Equity curve from trade returns
cumret = 1.0
equity_vals = np.full(n, float(capital))
trade_idx = 0
in_trade = False
t_entry_price = None
t_dir = None
for i in range(n):
s = sig_arr[i]
c = price_arr[i]
if s != 0.0 and s != t_dir:
if t_dir is not None and t_entry_price is not None and t_entry_price != 0:
t_ret = t_dir * (c - t_entry_price) / t_entry_price - cost
cumret *= (1 + t_ret)
t_entry_price = c
t_dir = s
equity_vals[i] = capital * cumret
# Bar returns for Sharpe
bar_returns = np.zeros(n)
for i in range(1, n):
if price_arr[i - 1] != 0 and last_dir is not None:
bar_returns[i] = sig_arr[i - 1] * (price_arr[i] - price_arr[i - 1]) / price_arr[i - 1] if sig_arr[i - 1] != 0 else 0.0
return {
"equity": pd.Series(equity_vals, index=close.index),
"trade_returns": trade_returns,
"long_returns": long_returns,
"short_returns": short_returns,
"bar_returns": bar_returns,
"trade_log": trade_log,
}
def compute_trade_stats(trades, capital=10000):
"""Single source of truth for trade statistics.
Every display path reads from this — no recomputation anywhere.
All values are rounded and JSON-safe (no inf/nan).
"""
if not trades:
return {"n": 0, "wins": 0, "losses": 0, "wr": 0, "avg": 0,
"best": 0, "worst": 0, "ret": 0, "np": 0, "mdd": 0,
"pf": 0, "rr": 0, "expect": 0}
w = [r for r in trades if r > 0]
l = [r for r in trades if r < 0]
cumret = 1.0
for r in trades:
cumret *= (1 + r)
net_p = capital * (cumret - 1)
# Max drawdown
eq = np.cumprod([1.0] + [1 + r for r in trades])
peak = np.maximum.accumulate(eq)
mdd = float(((eq - peak) / peak).min()) if len(eq) > 1 else 0.0
# Profit Factor
gross_w = sum(w) if w else 0
gross_l = abs(sum(l)) if l else 0
pf = gross_w / gross_l if gross_l > 0 else (9999.0 if gross_w > 0 else 0)
# Risk:Reward
avg_w = float(np.mean(w)) if w else 0
avg_l = abs(float(np.mean(l))) if l else 0
rr = avg_w / avg_l if avg_l > 0 else (9999.0 if avg_w > 0 else 0)
# Expectancy
expect = net_p / len(trades)
return {
"n": len(trades), "wins": len(w), "losses": len(l),
"wr": round(len(w) / len(trades), 4),
"avg": round(float(np.mean(trades)), 6),
"best": round(max(w), 6) if w else 0,
"worst": round(min(l), 6) if l else 0,
"ret": round(cumret - 1, 6),
"np": round(net_p, 2),
"mdd": round(mdd, 6),
"pf": round(pf, 2),
"rr": round(rr, 2),
"expect": round(expect, 2),
}
def compute_metrics(bt_result, close_test, capital=10000):
"""Compute all standard metrics from backtest result.
Uses trade-level compounding (same as webapp _trade_stats) for accuracy.
Returns: dict with total_ret, bh_ret, sharpe_strat, sharpe_bh, mdd, n_trades
"""
equity = bt_result["equity"]
trade_returns = bt_result["trade_returns"]
# Total return — trade-level compounding (matches webapp)
if trade_returns:
cumret = 1.0
for r in trade_returns:
cumret *= (1 + r)
total_ret = cumret - 1
else:
total_ret = 0.0
# Buy and hold
bh_equity = capital * (close_test / close_test.iloc[0])
bh_ret = (bh_equity.iloc[-1] - capital) / capital if capital != 0 else 0.0
# Sharpe ratio — trade-level (matches webapp: sqrt(252*26) annualization)
if len(trade_returns) >= 2 and float(np.std(trade_returns)) > 0:
sharpe_strat = float(np.mean(trade_returns) / np.std(trade_returns) * np.sqrt(252 * 26))
else:
sharpe_strat = 0.0
bh_rets = bh_equity.pct_change().dropna()
if len(bh_rets) > 1 and bh_rets.std() != 0:
sharpe_bh = float((bh_rets.mean() / bh_rets.std()) * np.sqrt(252 * 24 * 4))
else:
sharpe_bh = 0.0
# Max drawdown — trade-level (matches webapp)
if trade_returns:
eq = np.cumprod([1.0] + [1 + r for r in trade_returns])
peak = np.maximum.accumulate(eq)
mdd = float(((eq - peak) / peak).min()) if len(eq) > 1 else 0.0
else:
mdd = 0.0
return {
"total_ret": float(total_ret),
"bh_ret": float(bh_ret),
"sharpe_strat": float(sharpe_strat) if not np.isnan(sharpe_strat) else 0.0,
"sharpe_bh": float(sharpe_bh) if not np.isnan(sharpe_bh) else 0.0,
"mdd": float(mdd),
"n_trades": len(trade_returns),
}
# Diagnostics line/histogram series (equity / drawdown / rolling_acc / conf_hist)
# only feed the small Diagnostics charts — they're never used by the price chart
# or scroll-back. On a 1-min model trained over the (2.2-capped) window these are
# still ~30k points each; downsample to a visually-identical resolution before the
# dict leaves the trainer so it doesn't carry that into Server-A RAM / Postgres.
_RESULTS_SERIES_MAX = 5000
def _downsample_idx(n, cap=_RESULTS_SERIES_MAX):
"""Evenly-spaced index list spanning [0, n-1] (first+last always kept), or
None when no downsampling is needed (n <= cap)."""
if n <= cap:
return None
return np.unique(np.linspace(0, n - 1, cap).astype(int)).tolist()
def _take(arr, idx):
"""Subset a list by an index list (idx may be None → return arr unchanged)."""
if idx is None or not isinstance(arr, list):
return arr
return [arr[i] for i in idx]
# trade_log / train_trade_log are lists of per-trade dicts (display-only — the
# Trade Log tab). They scale with TRADE count, not bar count, so the bar-window
# cap (Phase 2.2) doesn't bound them — a degenerate near-every-bar model can put
# 10k+ trade dicts in the blob (>3 MB). Cap each (independently — a small-N model
# keeps every trade) to the most-recent N, recording `*_total` + `*_truncated`
# so the true count is still reported. Real strategies have far fewer than
# _TRADE_LOG_MAX trades, so this only ever bites pathological models.
_TRADE_LOG_MAX = 5000
def _cap_trade_log(tl):
"""Return (capped_list, original_len, was_truncated)."""
if not isinstance(tl, list) or len(tl) <= _TRADE_LOG_MAX:
return tl, (len(tl) if isinstance(tl, list) else 0), False
return tl[-_TRADE_LOG_MAX:], len(tl), True
def build_return_dict(split_result, bt_result, metrics, model, feature_cols,
signal_full, p_pos_test, p_neg_test, custom_figs=None,
bt_train_result=None, pre_stats=None):
"""Assemble the complete return dict. Handles ALL serialization.
Never returns Timestamps, numpy arrays, or non-JSON types.
Returns: JSON-safe dict with all required keys
"""
df = split_result["df"]
close = split_result["close"]
close_test = split_result["close_test"]
X_test = split_result["X_test"]
y_test = split_result["y_test"]
equity = bt_result["equity"]
bar_returns = bt_result["bar_returns"]
# OHLC
ohlc_dates = [str(x) for x in df.index.tolist()]
def _safe_list(arr):
return [float(x) if (x is not None and not np.isnan(x) and not np.isinf(x)) else None
for x in arr]
# Overlays
bb, ma = compute_overlays(close, df.index)
# Buy and hold equity
capital = equity.iloc[0] if len(equity) > 0 else 10000
bh_equity = capital * (close_test / close_test.iloc[0])
# Confusion matrix
from sklearn.metrics import confusion_matrix
pred_test = model.predict(X_test)
y_test_arr = np.asarray(y_test)
cm = confusion_matrix(y_test_arr, pred_test, labels=[-1, 0, 1])
# Rolling accuracy
sig_arr = signal_full.reindex(close_test.index).values
correct = pd.Series((pred_test == y_test_arr).astype(float), index=X_test.index)
active_test = pd.Series(sig_arr != 0, index=close_test.index) if len(sig_arr) == len(close_test) else pd.Series(True, index=close_test.index)
correct_active = correct.where(active_test, other=np.nan)
rolling_acc = correct_active.rolling(30, min_periods=1).mean()
# Feature importance
importances = model.feature_importances_
fi_pairs = sorted(zip(feature_cols, importances), key=lambda x: x[1])[-15:]
# Drawdown
rolling_max = equity.cummax()
drawdown = (equity - rolling_max) / rolling_max.replace(0, np.nan)
drawdown = drawdown.fillna(0.0)
# ── Downsample the Diagnostics-only series (see _downsample_idx) ──────────
_eq_dates = [str(x) for x in close_test.index.tolist()]
_eq_strat = _safe_list(equity.values)
_eq_bh = _safe_list(bh_equity.values)
_eq_idx = _downsample_idx(len(_eq_dates))
_eq_dates, _eq_strat, _eq_bh = _take(_eq_dates, _eq_idx), _take(_eq_strat, _eq_idx), _take(_eq_bh, _eq_idx)
_ra_dates = [str(x) for x in rolling_acc.index.tolist()]
_ra_vals = [float(x) if (not np.isnan(x) and not np.isinf(x)) else None for x in rolling_acc.values]
_ra_idx = _downsample_idx(len(_ra_dates))
_ra_dates, _ra_vals = _take(_ra_dates, _ra_idx), _take(_ra_vals, _ra_idx)
_dd_dates = [str(x) for x in drawdown.index.tolist()]
_dd_vals = _safe_list(drawdown.values)
_dd_idx = _downsample_idx(len(_dd_dates))
_dd_dates, _dd_vals = _take(_dd_dates, _dd_idx), _take(_dd_vals, _dd_idx)
_cp_pos = [float(x) for x in (p_pos_test.tolist() if hasattr(p_pos_test, 'tolist') else list(p_pos_test))]
_cp_neg = [float(x) for x in (p_neg_test.tolist() if hasattr(p_neg_test, 'tolist') else list(p_neg_test))]
_cp_pos = _take(_cp_pos, _downsample_idx(len(_cp_pos)))
_cp_neg = _take(_cp_neg, _downsample_idx(len(_cp_neg)))
# ── Trade logs — display-only (Trade Log tab); cap to most-recent N with a
# `_total` field so the true count is still reported (see _cap_trade_log).
# NB: ret_dist arrays are left FULL — a downstream path in callbacks.py
# recomputes n_trades/win-rate from len(ret_dist), so a sample would skew
# the displayed counts; they're small anyway and gzip handles them.
_tl_test, _tl_test_n, _tl_test_tr = _cap_trade_log(bt_result.get("trade_log", []))
_tl_tr, _tl_tr_n, _tl_tr_tr = _cap_trade_log(bt_train_result.get("trade_log", []) if bt_train_result else [])
return {
"ohlc": {
"dates": ohlc_dates,
"open": _safe_list(df["open"].values),
"high": _safe_list(df["high"].values),
"low": _safe_list(df["low"].values),
"close": _safe_list(df["close"].values),
},
"signals": {
"dates": [str(x) for x in signal_full.index.tolist()],
"values": [float(x) for x in signal_full.values],
},
"bb": bb,
"ma": ma,
"equity": {
"dates": _eq_dates,
"strategy": _eq_strat,
"bh": _eq_bh,
},
"feature_importance": {
"names": [p[0] for p in fi_pairs],
"values": [float(p[1]) for p in fi_pairs],
},
"conf_matrix": cm.tolist(),
"conf_hist": {
"p_pos": _cp_pos,
"p_neg": _cp_neg,
},
"rolling_acc": {
"dates": _ra_dates,
"values": _ra_vals,
},
"drawdown": {
"dates": _dd_dates,
"values": _dd_vals,
},
"ret_dist": [float(x) for x in bt_result["trade_returns"]],
"ret_dist_long": [float(x) for x in bt_result["long_returns"]],
"ret_dist_short": [float(x) for x in bt_result["short_returns"]],
"train_ret_dist": [float(x) for x in bt_train_result["trade_returns"]] if bt_train_result else [],
"train_ret_dist_long": [float(x) for x in bt_train_result["long_returns"]] if bt_train_result else [],
"train_ret_dist_short": [float(x) for x in bt_train_result["short_returns"]] if bt_train_result else [],
"trade_log": _tl_test,
"train_trade_log": _tl_tr,
"trade_log_total": _tl_test_n,
"train_trade_log_total": _tl_tr_n,
"trade_log_truncated": _tl_test_tr,
"train_trade_log_truncated": _tl_tr_tr,
**(pre_stats or {}),
"metrics": metrics,
"split_dt": split_result["split_dt"],
"split_idx": int(split_result["split_idx"]),
"n_train": int(split_result["n_train"]),
"n_test": int(split_result["n_test"]),
"feature_cols": list(feature_cols),
"custom_figs": custom_figs or [],
}
# ════════════════════════════════════════════════════════════════════════════
# STRATEGY FRAMEWORK v2 — Config-driven architecture
# Claude writes feature_engineering() + strategy_config(). Framework does rest.
# ════════════════════════════════════════════════════════════════════════════
import importlib
_MODEL_REGISTRY = {
"XGBClassifier": ("xgboost", "XGBClassifier"),
"RandomForestClassifier": ("sklearn.ensemble", "RandomForestClassifier"),
"GradientBoostingClassifier": ("sklearn.ensemble", "GradientBoostingClassifier"),
"LogisticRegression": ("sklearn.linear_model", "LogisticRegression"),
"ExtraTreesClassifier": ("sklearn.ensemble", "ExtraTreesClassifier"),
"AdaBoostClassifier": ("sklearn.ensemble", "AdaBoostClassifier"),
}
def _build_model_from_config(config, X_train, y_train_enc):
"""Build, fit, and wrap a model from strategy_config dict."""
model_type = config.get("model_type", "RandomForestClassifier")
model_params = dict(config.get("model_params", {}))
if model_type not in _MODEL_REGISTRY:
raise ValueError(f"Unknown model_type '{model_type}'. Valid: {list(_MODEL_REGISTRY.keys())}")
module_path, class_name = _MODEL_REGISTRY[model_type]
mod = importlib.import_module(module_path)
cls = getattr(mod, class_name)
# XGBoost defaults
if class_name == "XGBClassifier":
model_params.setdefault("use_label_encoder", False)
model_params.setdefault("eval_metric", "mlogloss")
model_params.setdefault("tree_method", "hist")
# Determinism > speed (2026-05-25). XGBoost hist with n_jobs=-1 is
# NON-reproducible even with random_state set — the parallel histogram
# gradient-sum order varies across threads, so the SAME code + data
# gives a slightly different model (and backtest) every run. Forcing
# single-thread makes training bit-reproducible so: (a) a user who
# copies a strategy and reruns it gets identical numbers, (b) the
# community "Live" score matches a redeploy, (c) "same code, different
# result" support reports go away. Cost: single-threaded XGB (a few
# seconds slower on large windows; hist is fast so it's minor). FORCED
# (not setdefault) so the guarantee can't be silently broken by a
# strategy passing n_jobs. Exact reproducibility holds within the
# platform (pinned versions / same Modal image); a user's own machine
# with different xgboost/numpy/CPU can still differ in low-order bits.
model_params["n_jobs"] = 1
# Common defaults
model_params.setdefault("random_state", 42)
from model_wrapper import ModelWrapper
clf = cls(**model_params)
clf.fit(X_train, y_train_enc)
enc = LabelEncoder()
enc.fit([-1, 0, 1])
return ModelWrapper(clf, original_classes=enc.classes_, n_features=X_train.shape[1])
def _generate_signals(model, X, threshold):
"""Framework-owned signal generation. Deterministic threshold logic."""
proba = model.predict_proba(X)
classes = list(model.classes_)
idx_pos = classes.index(1) if 1 in classes else None
idx_neg = classes.index(-1) if -1 in classes else None
p_pos = proba[:, idx_pos] if idx_pos is not None else np.zeros(len(X))
p_neg = proba[:, idx_neg] if idx_neg is not None else np.zeros(len(X))
signal_vals = np.zeros(len(X))
signal_vals = np.where(p_pos >= threshold, 1.0, signal_vals)
signal_vals = np.where(p_neg >= threshold, -1.0, signal_vals)
# Both exceed: pick stronger
both = (p_pos >= threshold) & (p_neg >= threshold)
signal_vals[both] = np.where(p_pos[both] >= p_neg[both], 1.0, -1.0)
return pd.Series(signal_vals, index=X.index), p_pos, p_neg
# ── Filter functions (all no-ops when config value is None) ──────────────
def _apply_direction_filter(signal, direction):
"""Zero out signals that don't match allowed direction."""
if direction is None or direction == "both":
return signal
s = signal.copy()
if direction == "long":
s[s < 0] = 0.0
elif direction == "short":
s[s > 0] = 0.0
return s
def _apply_session_filter(signal, index, session_hours):
"""Zero out signals outside session hours [start, end] UTC."""
if session_hours is None:
return signal
s = signal.copy()
start_h, end_h = session_hours[0], session_hours[1]
hours = index.hour
if start_h <= end_h:
mask = (hours >= start_h) & (hours < end_h)
else: # wrap around midnight, e.g. [22, 6]
mask = (hours >= start_h) | (hours < end_h)
s[~mask] = 0.0
return s
def _apply_atr_filter(signal, close, high, low, min_atr):
"""Zero out signals when NATR(14) is below threshold."""
if min_atr is None:
return signal
hl = high - low
hc = (high - close.shift(1)).abs()
lc = (low - close.shift(1)).abs()
tr = pd.concat([hl, hc, lc], axis=1).max(axis=1)
atr14 = tr.ewm(com=13, adjust=False).mean()
natr = atr14 / close.replace(0, np.nan)
s = signal.copy()
s[natr < min_atr] = 0.0
return s
def _apply_trend_filter(signal, close, trend_filter):
"""Only allow signals aligned with trend. e.g. 'sma_50': longs above SMA, shorts below."""
if trend_filter is None:
return signal
# Parse: "sma_50" → SMA with period 50
parts = trend_filter.lower().replace("-", "_").split("_")
if len(parts) >= 2 and parts[0] in ("sma", "ema"):
period = int(parts[1])
else:
return signal # unknown filter, skip
if parts[0] == "sma":
trend_line = close.rolling(period).mean()
else:
trend_line = close.ewm(span=period, adjust=False).mean()
s = signal.copy()
# Longs only above trend, shorts only below
s[(s > 0) & (close < trend_line)] = 0.0
s[(s < 0) & (close > trend_line)] = 0.0
return s
# ── run_backtest_v2: framework-owned SL/TP/cooldown/position management ──
def run_backtest_v2(signal, close, high, low, config, capital=10000, cost=2e-5):
"""Backtest with SL/TP/cooldown/direction handling built into the engine.
Unlike run_backtest (v1), this function handles position exits internally.
Returns: same dict shape as run_backtest()
"""
stop_loss = config.get("stop_loss")
take_profit = config.get("take_profit")
cooldown = config.get("cooldown", 0)
on_opposite = config.get("on_opposite", "reverse")
sig_arr = signal.values
close_arr = close.values
high_arr = high.values
low_arr = low.values
idx = signal.index
n = len(close_arr)
trade_returns = []
long_returns = []
short_returns = []
trade_log = []
equity_vals = np.full(n, float(capital))
cumret = 1.0
position = 0.0 # current direction: 1.0, -1.0, or 0.0 (flat)
entry_price = None
entry_bar = None # index into arrays for entry time
cooldown_remaining = 0
def _log_trade(exit_bar, exit_px, ret, reason):
trade_log.append({
"type": "Buy" if position == 1.0 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[exit_bar]),
"entry_price": round(entry_price, 5),
"exit_price": round(exit_px, 5),
"pnl": round(position * (exit_px - entry_price), 5),
"pnl_pct": round(ret * 100, 3),
"exit_reason": reason,
})
for i in range(n):
c = close_arr[i]
h = high_arr[i]
lo = low_arr[i]
s = sig_arr[i]
# 1. Check SL/TP if in trade
if position != 0.0 and entry_price is not None:
hit_sl = False
hit_tp = False
exit_price = None
if position == 1.0: # long
if stop_loss is not None and lo <= entry_price * (1 - stop_loss):
hit_sl = True
exit_price = entry_price * (1 - stop_loss)
elif take_profit is not None and h >= entry_price * (1 + take_profit):
hit_tp = True
exit_price = entry_price * (1 + take_profit)
else: # short
if stop_loss is not None and h >= entry_price * (1 + stop_loss):
hit_sl = True
exit_price = entry_price * (1 + stop_loss)
elif take_profit is not None and lo <= entry_price * (1 - take_profit):
hit_tp = True
exit_price = entry_price * (1 - take_profit)
if hit_sl or hit_tp:
ret = float(position * (exit_price - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, exit_price, ret, "SL" if hit_sl else "TP")
cumret *= (1 + ret)
position = 0.0
entry_price = None
entry_bar = None
cooldown_remaining = cooldown
equity_vals[i] = capital * cumret
continue
# 2. Cooldown
if cooldown_remaining > 0:
cooldown_remaining -= 1
equity_vals[i] = capital * cumret
continue
# 3. Signal processing
if s != 0.0:
if position == 0.0:
# Open new trade
position = s
entry_price = c
entry_bar = i
elif s != position:
# Opposite signal
if on_opposite == "reverse":
# Close current + open opposite
ret = float(position * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, c, ret, "signal")
cumret *= (1 + ret)
position = s
entry_price = c
entry_bar = i
else: # close_only
# Close current, go flat
ret = float(position * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, c, ret, "close_only")
cumret *= (1 + ret)
position = 0.0
entry_price = None
entry_bar = None
cooldown_remaining = cooldown
equity_vals[i] = capital * cumret
# Close last open trade at final close
if position != 0.0 and entry_price is not None and n > 0 and entry_price != 0:
c = close_arr[-1]
ret = float(position * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(n - 1, c, ret, "end")
cumret *= (1 + ret)
equity_vals[-1] = capital * cumret
# Bar returns for Sharpe (approximate)
bar_returns = np.zeros(n)
for i in range(1, n):
if close_arr[i - 1] != 0 and sig_arr[i - 1] != 0:
bar_returns[i] = sig_arr[i - 1] * (close_arr[i] - close_arr[i - 1]) / close_arr[i - 1]
return {
"equity": pd.Series(equity_vals, index=close.index),
"trade_returns": trade_returns,
"long_returns": long_returns,
"short_returns": short_returns,
"bar_returns": bar_returns,
"trade_log": trade_log,
}
# ── run_strategy: the v2 orchestrator ────────────────────────────────────
def run_strategy(feature_fn, config_fn, data_path, start_date="", end_date="",
validation_date="", train_split=0.7, register_model_fn=None):
"""Config-driven strategy execution. Claude writes feature_fn + config_fn,
framework does everything else.
Returns: results dict (same format as webapp expects)
"""
config = config_fn()
# Auto-correct SL/TP if Claude passed percentage instead of decimal
for _key in ("stop_loss", "take_profit"):
_val = config.get(_key)
if _val is not None and _val > 0.1: # >10% is almost certainly a percentage
config[_key] = _val / 100.0
print(f"[strategy] Auto-corrected {_key}: {_val} -> {config[_key]} (was percentage, converted to decimal)")
# 1. Load data
df, close, open_, high, low = load_ohlc(data_path, start_date, end_date)
# 2. Feature engineering (Claude's function)
df = feature_fn(df, close, open_, high, low)
close = df["close"]
open_ = df["open"]
high = df["high"]
low = df["low"]
# 3. Warm-up detection: drop rows where features have NaN BEFORE any fill
feature_cols = [c for c in df.columns if c not in ("open", "high", "low", "close")]
raw_nans = df[feature_cols].isna().any(axis=1)
valid_rows = ~raw_nans
if valid_rows.any():
first_valid = valid_rows.idxmax()
if raw_nans.loc[:first_valid].any():
df = df.loc[first_valid:].copy()
close = df["close"]
open_ = df["open"]
high = df["high"]
low = df["low"]
# 4. Target
horizon = config.get("target_horizon", 4)
target = make_target(close, horizon=horizon)
# 5. Split (ffill only within each partition — no bfill leak)
mask = target.notna()
df = df[mask].copy()
target = target[mask]
close = df["close"]
high = df["high"]
low = df["low"]
X = df[feature_cols].copy()
X = X.replace([np.inf, -np.inf], np.nan)
if validation_date:
split_idx = len(df[df.index <= validation_date])
else:
split_idx = int(len(df) * train_split)
split_idx = max(1, min(split_idx, len(df) - 1))
# ffill within train and test separately (no leak)
X_train = X.iloc[:split_idx].ffill().fillna(0.0)
X_test = X.iloc[split_idx:].ffill().fillna(0.0)
X = pd.concat([X_train, X_test])
y_train = target.iloc[:split_idx]
y_test = target.iloc[split_idx:]
close_train = close.iloc[:split_idx]
close_test = close.iloc[split_idx:]
high_test = high.iloc[split_idx:]
low_test = low.iloc[split_idx:]
enc = LabelEncoder()
enc.fit([-1, 0, 1])
y_train_enc = enc.transform(y_train)
y_test_enc = enc.transform(y_test)
split_dt = str(df.index[split_idx])
sp = {
"df": df, "X_train": X_train, "X_test": X_test,
"y_train": y_train, "y_test": y_test,
"y_train_enc": y_train_enc, "y_test_enc": y_test_enc,
"enc": enc,
"close": close, "close_train": close_train, "close_test": close_test,
"split_idx": split_idx, "split_dt": split_dt,
"n_train": len(X_train), "n_test": len(X_test),
}
# 6. Build model from config
model = _build_model_from_config(config, X_train, y_train_enc)
# 7. Generate signals
threshold = config.get("signal_threshold", 0.55)
signal_train, p_pos_train, p_neg_train = _generate_signals(model, X_train, threshold)
signal_test, p_pos_test, p_neg_test = _generate_signals(model, X_test, threshold)
# 8. Apply filters (order: direction → session → ATR → trend)
direction = config.get("direction", "both")
signal_test = _apply_direction_filter(signal_test, direction)
signal_train = _apply_direction_filter(signal_train, direction)
session_filter = config.get("session_filter")
signal_test = _apply_session_filter(signal_test, signal_test.index, session_filter)
signal_train = _apply_session_filter(signal_train, signal_train.index, session_filter)
min_atr = config.get("min_atr")
if min_atr is not None:
signal_test = _apply_atr_filter(signal_test, close_test, high_test, low_test, min_atr)
trend_filter = config.get("trend_filter")
if trend_filter is not None:
signal_test = _apply_trend_filter(signal_test, close_test, trend_filter)
signal_full = pd.concat([signal_train, signal_test])
# 9. Backtest with SL/TP/cooldown (test + train)
high_train = high.iloc[:split_idx]
low_train = low.iloc[:split_idx]
has_risk = (config.get("stop_loss") is not None or
config.get("take_profit") is not None or
config.get("cooldown", 0) > 0 or
config.get("on_opposite", "reverse") != "reverse")
if has_risk:
bt = run_backtest_v2(signal_test, close_test, high_test, low_test, config, capital=10000)
bt_train = run_backtest_v2(signal_train, close_train, high_train, low_train, config, capital=10000)
else:
bt = run_backtest(signal_test, close_test, capital=10000)
bt_train = run_backtest(signal_train, close_train, capital=10000)
# 10. Metrics
metrics = compute_metrics(bt, close_test, capital=10000)
# 11. Pre-compute all trade stats (single source of truth)
pre_stats = {
"train_stats": compute_trade_stats(bt_train.get("trade_returns", []), capital=10000),
"test_stats": compute_trade_stats(bt.get("trade_returns", []), capital=10000),
"long_stats": compute_trade_stats(bt.get("long_returns", []), capital=10000),
"short_stats": compute_trade_stats(bt.get("short_returns", []), capital=10000),
}
# 12. Register model
if register_model_fn is not None:
register_model_fn(model)
# 13. Build return dict
return build_return_dict(sp, bt, metrics, model, feature_cols,
signal_full, p_pos_test, p_neg_test, custom_figs=[],
bt_train_result=bt_train, pre_stats=pre_stats)
# ── End strategy_utils ──
DATA_PATH = '/root/Desktop/QuantifyMe/data/ohlc/GBPUSD_15min.parquet'
START_DATE = '2026-04-15'
END_DATE = '2026-05-25'
VALIDATION_DATE = ""
TRAIN_SPLIT = 0.7
# SECTION 1 — FEATURE ENGINEERING
def feature_engineering(df, close, open_, high, low):
# ── EMA crossover core signals ──────────────────────────────────────────
ema9 = close.ewm(span=9, adjust=False).mean()
ema21 = close.ewm(span=21, adjust=False).mean()
ema50 = close.ewm(span=50, adjust=False).mean()
ema200 = close.ewm(span=200, adjust=False).mean()
df["ema9"] = ema9
df["ema21"] = ema21
df["ema50"] = ema50
df["ema200"] = ema200
# Raw spread and normalised spread
df["ema_diff"] = ema9 - ema21
df["ema_diff_norm"] = (ema9 - ema21) / close
# Cross signal: +1 when ema9 > ema21, -1 otherwise
df["ema_cross_sign"] = np.where(ema9 > ema21, 1.0, -1.0)
# Momentum of the spread (rate of change of spread)
df["ema_diff_roc1"] = df["ema_diff"].diff(1)
df["ema_diff_roc3"] = df["ema_diff"].diff(3)
# Distance of price from ema50 and ema200 (normalised)
df["dist_ema50"] = (close - ema50) / close
df["dist_ema200"] = (close - ema200) / close
# ── RSI (14) ────────────────────────────────────────────────────────────
delta = close.diff()
gain = delta.clip(lower=0.0)
loss = (-delta).clip(lower=0.0)
avg_g = gain.ewm(com=13, adjust=False).mean()
avg_l = loss.ewm(com=13, adjust=False).mean()
rs = avg_g / avg_l.replace(0.0, np.nan)
rsi14 = 100.0 - 100.0 / (1.0 + rs)
df["rsi14"] = rsi14
# RSI normalised and centred
df["rsi14_norm"] = (rsi14 - 50.0) / 50.0
# ── MACD ────────────────────────────────────────────────────────────────
ema12 = close.ewm(span=12, adjust=False).mean()
ema26 = close.ewm(span=26, adjust=False).mean()
macd_line = ema12 - ema26
signal_ln = macd_line.ewm(span=9, adjust=False).mean()
macd_hist = macd_line - signal_ln
df["macd_line"] = macd_line / close
df["macd_signal"] = signal_ln / close
df["macd_hist"] = macd_hist / close
df["macd_cross"] = np.where(macd_line > signal_ln, 1.0, -1.0)
# ── Bollinger Bands (20, 2) ──────────────────────────────────────────────
bb_mid = close.rolling(20).mean()
bb_std = close.rolling(20).std(ddof=0)
bb_upper = bb_mid + 2.0 * bb_std
bb_lower = bb_mid - 2.0 * bb_std
bb_width = (bb_upper - bb_lower) / bb_mid.replace(0.0, np.nan)
bb_pct = (close - bb_lower) / (bb_upper - bb_lower).replace(0.0, np.nan)
df["bb_width"] = bb_width
df["bb_pct"] = bb_pct
# ── ATR (14) ─────────────────────────────────────────────────────────────
tr = pd.concat([
high - low,
(high - close.shift(1)).abs(),
(low - close.shift(1)).abs()
], axis=1).max(axis=1)
atr14 = tr.ewm(com=13, adjust=False).mean()
df["atr14"] = atr14
df["natr14"] = atr14 / close # normalised ATR (volatility proxy)
# ── Stochastic %K / %D (14, 3) ──────────────────────────────────────────
low14 = low.rolling(14).min()
high14 = high.rolling(14).max()
stoch_k = 100.0 * (close - low14) / (high14 - low14).replace(0.0, np.nan)
stoch_d = stoch_k.rolling(3).mean()
df["stoch_k"] = stoch_k / 100.0
df["stoch_d"] = stoch_d / 100.0
df["stoch_diff"] = (stoch_k - stoch_d) / 100.0
# ── Rate of Change ───────────────────────────────────────────────────────
df["roc1"] = close.pct_change(1)
df["roc4"] = close.pct_change(4)
df["roc8"] = close.pct_change(8)
df["roc16"] = close.pct_change(16)
# ── Candle features ──────────────────────────────────────────────────────
body = (close - open_).abs()
candle_rng = (high - low).replace(0.0, np.nan)
df["body_ratio"] = body / candle_rng
df["upper_wick"] = (high - pd.concat([close, open_], axis=1).max(axis=1)) / candle_rng
df["lower_wick"] = (pd.concat([close, open_], axis=1).min(axis=1) - low) / candle_rng
df["candle_dir"] = np.where(close >= open_, 1.0, -1.0)
# ── Volume-like proxy: range relative to rolling average ────────────────
df["range_ratio"] = candle_rng / candle_rng.rolling(20).mean()
# ── Lagged EMA diff features ─────────────────────────────────────────────
for lag in [1, 2, 3, 4]:
df[f"ema_diff_lag{lag}"] = df["ema_diff_norm"].shift(lag)
# ── Lagged RSI ───────────────────────────────────────────────────────────
for lag in [1, 2, 4]:
df[f"rsi14_lag{lag}"] = df["rsi14_norm"].shift(lag)
# ── Rolling volatility (std of returns) ──────────────────────────────────
ret = close.pct_change()
df["vol_8"] = ret.rolling(8).std()
df["vol_16"] = ret.rolling(16).std()
df["vol_32"] = ret.rolling(32).std()
# ── Trend strength: ADX-like (simplified) ────────────────────────────────
plus_dm = (high.diff()).clip(lower=0.0)
minus_dm = (-low.diff()).clip(lower=0.0)
overlap = pd.concat([plus_dm, minus_dm], axis=1).min(axis=1)
plus_dm = plus_dm - overlap
minus_dm = minus_dm - overlap
smooth_tr = tr.ewm(com=13, adjust=False).mean()
plus_di = 100.0 * plus_dm.ewm(com=13, adjust=False).mean() / smooth_tr.replace(0.0, np.nan)
minus_di = 100.0 * minus_dm.ewm(com=13, adjust=False).mean() / smooth_tr.replace(0.0, np.nan)
di_sum = (plus_di + minus_di).replace(0.0, np.nan)
adx = ((plus_di - minus_di).abs() / di_sum * 100.0).ewm(com=13, adjust=False).mean()
df["adx"] = adx / 100.0
df["plus_di"] = plus_di / 100.0
df["minus_di"] = minus_di / 100.0
# ── Session hour (UTC) ───────────────────────────────────────────────────
if hasattr(df.index, "hour"):
df["hour_sin"] = np.sin(2.0 * np.pi * df.index.hour / 24.0)
df["hour_cos"] = np.cos(2.0 * np.pi * df.index.hour / 24.0)
else:
df["hour_sin"] = 0.0
df["hour_cos"] = 1.0
# ── Fill any NaN from warm-up periods ────────────────────────────────────
df = df.bfill().ffill()
return df
# SECTION 2 — STRATEGY CONFIG
def strategy_config():
return {
"title": "EMA 9/21 Crossover + MACD Momentum (XGBoost)",
"model_type": "XGBClassifier",
"model_params": {
"n_estimators": 400,
"max_depth": 4,
"learning_rate": 0.04,
"subsample": 0.80,
"colsample_bytree": 0.75,
"min_child_weight": 3,
"gamma": 0.10,
"reg_alpha": 0.05,
"reg_lambda": 1.50,
"objective": "binary:logistic",
"tree_method": "hist",
"random_state": 42,
},
"signal_threshold": 0.55,
"direction": "both",
"stop_loss": 0.0030,
"take_profit": 0.0060,
"cooldown": 0,
"max_positions": 1,
"on_opposite": "reverse",
"session_filter": [6, 20],
"min_atr": None,
"trend_filter": "sma_50",
"target_horizon": 4,
"objective": (
"Maximize Sharpe ratio on EUR/USD 15-min data. "
"Core signal: EMA(9) vs EMA(21) crossover enriched with MACD, RSI, "
"Bollinger %B, Stochastic, ATR, ADX, candle structure and rolling "
"volatility. XGBoost with moderate depth (4) and strong regularisation "
"(gamma, alpha, lambda) prevents overfitting on ~6 weeks of intraday data. "
"A 0.55 probability threshold filters low-confidence signals. "
"A 2:1 TP:SL ratio (30 bp SL / 60 bp TP) improves the reward-risk "
"balance. Session filter [6,20] UTC keeps the model away from the thin "
"Asian pre-open. trend_filter sma_50 aligns entries with the prevailing "
"short-term trend to reduce chop. Cooldown=0 and reverse-on-opposite "
"allow continuous participation in trending EMA crossover moves."
),
"notes": (
"round-trip cost 2e-5 is accounted for by the framework. "
"target_horizon=4 bars (1 hour ahead) suits EMA crossover which "
"generates medium-frequency signals rather than tick-level scalps. "
"All features are normalised or expressed as ratios to minimise "
"scale sensitivity for the logistic-objective XGBoost."
),
}
# ── Framework v2: auto-generated wrapper ──
def train_and_backtest():
_vd = VALIDATION_DATE if 'VALIDATION_DATE' in globals() else ''
_ts = TRAIN_SPLIT if 'TRAIN_SPLIT' in globals() else 0.7
return run_strategy(
feature_engineering, strategy_config,
DATA_PATH, START_DATE, END_DATE,
_vd, _ts,
register_model_fn=register_model
)
|
||||||||||
|
—
|
RSI mean-reversion
|
M
@malcolmtan
|
EURUSD | 1min | 65.2%68.2% | +0.01%-0.25% | 1.010.99 | 0.82%0.82% | 2344 |
|
# ╔══════════════════════════════════════════════════════════════╗
# ║ STRATEGY REQUEST LOG ║
# ╚══════════════════════════════════════════════════════════════╝
# Generated : 2026-05-25 02:26:36
# Model : XGBoost
# Feature Eng. : buy when RSI(14) crosses up from below 30, sell when it crosses down from above 70 + Auto-add features: ON
# Signal / Entry : —
# Optimization : —
# Risk Mgmt : —
# Risk Filter : —
# ══════════════════════════════════════════════════════════════
# ============================================================
# SECTION 0 — IMPORTS & CONSTANTS
import numpy as np
import pandas as pd
# ── Inlined strategy_utils ──
"""
strategy_utils.py — Standard utility functions for generated strategies.
Claude imports these instead of writing boilerplate from scratch.
This ensures consistent behavior across all generated strategies.
"""
import numpy as np
import pandas as pd
from sklearn.preprocessing import LabelEncoder
# Max backtest window per timeframe. A finer timeframe over a longer window
# blows up the results dict / parquet load / Modal train time (the 2026-05-12
# OOM was a 1-min × multi-year sweep) — and a 1-min strategy gains nothing from
# 2 years of 1-min bars. Enforced HERE because every training path (UI / API /
# Modal) funnels through run_strategy → load_ohlc. Env-overridable so a future
# "max plan" / dedicated-server tier can lift it.
_TF_MAX_DAYS = {
"1min": 30,
"5min": 90,
"15min": 365,
"1h": 730,
}
def _fetch_ohlc_from_internal(symbol: str, tf: str, start: str, end: str):
"""Phase 3.2: fetch parquet bytes from Server A's /internal/ohlc endpoint
instead of reading a local file. Used inside Modal containers / Mac worker
pool (Phase 3.4) so every train sees the same source of truth as the chart.
Returns: pd.DataFrame (parquet decoded), or raises on any failure so the
caller can fall back / surface a clear error in the job.
"""
import hashlib as _hashlib, hmac as _hmac, io as _io, os as _os
import urllib.request as _ur, urllib.parse as _urp
base = (_os.environ.get("QM_INTERNAL_OHLC_BASE") or "").rstrip("/")
secret = (_os.environ.get("INTERNAL_WS_SECRET") or "").strip()
if not base:
raise RuntimeError("QM_INTERNAL_OHLC_BASE not set")
if not secret:
raise RuntimeError("INTERNAL_WS_SECRET not set")
msg = f"{symbol}|{tf}|{start}|{end}".encode("utf-8")
sig = _hmac.new(secret.encode("utf-8"), msg, _hashlib.sha256).hexdigest()
qs = _urp.urlencode({
"symbol": symbol, "tf": tf,
"start": start, "end": end, "sig": sig,
})
url = f"{base}/internal/ohlc?{qs}"
req = _ur.Request(url, headers={"User-Agent": "qm-worker/1.0"})
with _ur.urlopen(req, timeout=30) as resp:
if resp.status != 200:
raise RuntimeError(f"/internal/ohlc returned {resp.status}")
payload = resp.read()
print(f"[load_ohlc:internal] {symbol} {tf} fetched {len(payload)} bytes", flush=True)
return pd.read_parquet(_io.BytesIO(payload))
def _parse_symbol_tf_from_path(data_path: str):
"""Pull SYMBOL + TF out of a path like .../EURUSD_1min.parquet."""
import os as _os, re as _re
base = _os.path.basename(str(data_path))
m = _re.match(r"^([A-Z]{6})_(\d+min|\d+h)\.parquet$", base)
if not m:
return None, None
return m.group(1), m.group(2)
def load_ohlc(data_path, start_date="", end_date=""):
"""Load OHLC parquet, sort index, filter dates. Always returns consistent format.
The lower bound is clamped per timeframe (see _TF_MAX_DAYS) — a request for
more history than the cap silently starts later.
Phase 3.2: when env QM_USE_INTERNAL_OHLC=="1", fetch over HTTP from
Server A's /internal/ohlc endpoint instead of pd.read_parquet on a local
file (which on Modal is a stale Volume snapshot). The endpoint applies the
same day-cap, so the local cap-check below is a defensive no-op in that
path. Flag defaults to "0" → unchanged behavior.
Returns: (df, close, open_, high, low)
"""
import os as _os, re as _re
_use_internal = _os.environ.get("QM_USE_INTERNAL_OHLC", "0") == "1"
if _use_internal:
_sym, _tf = _parse_symbol_tf_from_path(data_path)
if not _sym or not _tf:
raise RuntimeError(
f"QM_USE_INTERNAL_OHLC=1 but DATA_PATH basename does not match "
f"SYMBOL_TF.parquet: {data_path}"
)
df = _fetch_ohlc_from_internal(_sym, _tf, start_date or "", end_date or "")
else:
df = pd.read_parquet(data_path)
df.index = pd.to_datetime(df.index)
df = df.sort_index()
# Per-timeframe window cap (timeframe inferred from the parquet filename).
_m = _re.search(r"_(\d+min|\d+h)\.parquet$", _os.path.basename(str(data_path)))
_tf = _m.group(1) if _m else None
_max_days = _TF_MAX_DAYS.get(_tf)
if _max_days and _max_days > 0 and len(df):
_env_override = _os.environ.get(f"QM_MAX_DAYS_{_tf.upper()}")
if _env_override and _env_override.isdigit():
_max_days = int(_env_override)
try:
_eff_end = pd.Timestamp(end_date) if end_date else df.index.max()
_eff_end = min(_eff_end, df.index.max())
_floor = _eff_end - pd.Timedelta(days=_max_days)
_req_start = pd.Timestamp(start_date) if start_date else df.index.min()
if _req_start < _floor:
print(f"[load_ohlc] {_tf} backtest window capped to {_max_days}d: "
f"start {_req_start.date()} -> {_floor.date()}", flush=True)
start_date = _floor
except Exception as _e:
print(f"[load_ohlc] window-cap check skipped ({_e})", flush=True)
if start_date:
df = df[df.index >= start_date]
if end_date:
df = df[df.index <= end_date]
return df, df["close"], df["open"], df["high"], df["low"]
def make_target(close, horizon=4):
"""Create target: direction N bars ahead. Default 4 bars = 1 hour on 15-min data.
Returns: target (pd.Series of -1, 0, 1)
"""
return np.sign(close.shift(-horizon) - close)
def split_data(df, target, feature_cols, train_split=0.7, validation_date=""):
"""Train/test split. Handles both ratio and date-based splits.
Drops NaN from target before splitting. Encodes labels to [0,1,2].
Returns: dict with keys:
X_train, X_test, y_train, y_test,
y_train_enc, y_test_enc, enc,
close_train, close_test,
split_idx, split_dt, n_train, n_test
"""
# Drop NaN from target
mask = target.notna()
df = df[mask].copy()
target = target[mask]
close = df["close"]
# Build feature matrix
X = df[feature_cols].copy()
X = X.bfill().ffill()
X = X.replace([np.inf, -np.inf], np.nan).fillna(0.0)
# Split
if validation_date:
split_idx = len(df[df.index <= validation_date])
else:
split_idx = int(len(df) * train_split)
split_idx = max(1, min(split_idx, len(df) - 1))
X_train = X.iloc[:split_idx]
X_test = X.iloc[split_idx:]
y_train = target.iloc[:split_idx]
y_test = target.iloc[split_idx:]
close_train = close.iloc[:split_idx]
close_test = close.iloc[split_idx:]
split_dt = str(df.index[split_idx])
# Label encoding — always fit on [-1, 0, 1]
enc = LabelEncoder()
enc.fit([-1, 0, 1])
y_train_enc = enc.transform(y_train)
y_test_enc = enc.transform(y_test)
return {
"df": df, "X_train": X_train, "X_test": X_test,
"y_train": y_train, "y_test": y_test,
"y_train_enc": y_train_enc, "y_test_enc": y_test_enc,
"enc": enc,
"close": close, "close_train": close_train, "close_test": close_test,
"split_idx": split_idx, "split_dt": split_dt,
"n_train": len(X_train), "n_test": len(X_test),
}
def compute_overlays(close, df_index):
"""Compute BB and MA overlays on full dataset. Always consistent.
Returns: (bb_dict, ma_dict)
"""
bb_mid = close.rolling(20).mean()
bb_std = close.rolling(20).std()
bb_upper = bb_mid + 2 * bb_std
bb_lower = bb_mid - 2 * bb_std
ma50 = close.rolling(50).mean()
ma100 = close.rolling(100).mean()
ma200 = close.rolling(200).mean()
def _safe(s):
s = s.reindex(df_index).bfill().ffill()
return [float(x) if (x is not None and not np.isnan(x) and not np.isinf(x)) else None
for x in s.values]
bb = {"upper": _safe(bb_upper), "mid": _safe(bb_mid), "lower": _safe(bb_lower)}
ma = {"ma50": _safe(ma50), "ma100": _safe(ma100), "ma200": _safe(ma200)}
return bb, ma
def run_backtest(signal, close, capital=10000, cost=2e-5):
"""Run backtest with transaction costs.
Uses price-based trade returns (same as webapp _compute_trades).
Signal 0 = hold (keep current position), not close.
Returns: dict with equity, trade_returns, long_returns, short_returns, bar_returns
"""
sig_arr = signal.values
price_arr = close.values
idx = signal.index
n = len(price_arr)
# Trade returns — price-based (matches webapp _compute_trades exactly)
trade_returns = []
long_returns = []
short_returns = []
trade_log = []
last_dir = None
entry_price = None
entry_bar = None
for i in range(n):
s = sig_arr[i]
c = price_arr[i]
if s != 0.0 and s != last_dir:
# Direction change — close previous trade, open new
if last_dir is not None and entry_price is not None and entry_price != 0:
ret = float(last_dir * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if last_dir == 1:
long_returns.append(ret)
else:
short_returns.append(ret)
trade_log.append({
"type": "Buy" if last_dir == 1 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[i]),
"entry_price": round(entry_price, 5),
"exit_price": round(c, 5),
"pnl": round(last_dir * (c - entry_price), 5),
"pnl_pct": round(ret * 100, 3),
"exit_reason": "signal",
})
entry_price = c
entry_bar = i
last_dir = s
# Close last open trade
if last_dir is not None and entry_price is not None and n > 0 and entry_price != 0:
c = price_arr[-1]
ret = float(last_dir * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if last_dir == 1:
long_returns.append(ret)
else:
short_returns.append(ret)
trade_log.append({
"type": "Buy" if last_dir == 1 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[-1]),
"entry_price": round(entry_price, 5),
"exit_price": round(c, 5),
"pnl": round(last_dir * (c - entry_price), 5),
"pnl_pct": round(ret * 100, 3),
"exit_reason": "end",
})
# Equity curve from trade returns
cumret = 1.0
equity_vals = np.full(n, float(capital))
trade_idx = 0
in_trade = False
t_entry_price = None
t_dir = None
for i in range(n):
s = sig_arr[i]
c = price_arr[i]
if s != 0.0 and s != t_dir:
if t_dir is not None and t_entry_price is not None and t_entry_price != 0:
t_ret = t_dir * (c - t_entry_price) / t_entry_price - cost
cumret *= (1 + t_ret)
t_entry_price = c
t_dir = s
equity_vals[i] = capital * cumret
# Bar returns for Sharpe
bar_returns = np.zeros(n)
for i in range(1, n):
if price_arr[i - 1] != 0 and last_dir is not None:
bar_returns[i] = sig_arr[i - 1] * (price_arr[i] - price_arr[i - 1]) / price_arr[i - 1] if sig_arr[i - 1] != 0 else 0.0
return {
"equity": pd.Series(equity_vals, index=close.index),
"trade_returns": trade_returns,
"long_returns": long_returns,
"short_returns": short_returns,
"bar_returns": bar_returns,
"trade_log": trade_log,
}
def compute_trade_stats(trades, capital=10000):
"""Single source of truth for trade statistics.
Every display path reads from this — no recomputation anywhere.
All values are rounded and JSON-safe (no inf/nan).
"""
if not trades:
return {"n": 0, "wins": 0, "losses": 0, "wr": 0, "avg": 0,
"best": 0, "worst": 0, "ret": 0, "np": 0, "mdd": 0,
"pf": 0, "rr": 0, "expect": 0}
w = [r for r in trades if r > 0]
l = [r for r in trades if r < 0]
cumret = 1.0
for r in trades:
cumret *= (1 + r)
net_p = capital * (cumret - 1)
# Max drawdown
eq = np.cumprod([1.0] + [1 + r for r in trades])
peak = np.maximum.accumulate(eq)
mdd = float(((eq - peak) / peak).min()) if len(eq) > 1 else 0.0
# Profit Factor
gross_w = sum(w) if w else 0
gross_l = abs(sum(l)) if l else 0
pf = gross_w / gross_l if gross_l > 0 else (9999.0 if gross_w > 0 else 0)
# Risk:Reward
avg_w = float(np.mean(w)) if w else 0
avg_l = abs(float(np.mean(l))) if l else 0
rr = avg_w / avg_l if avg_l > 0 else (9999.0 if avg_w > 0 else 0)
# Expectancy
expect = net_p / len(trades)
return {
"n": len(trades), "wins": len(w), "losses": len(l),
"wr": round(len(w) / len(trades), 4),
"avg": round(float(np.mean(trades)), 6),
"best": round(max(w), 6) if w else 0,
"worst": round(min(l), 6) if l else 0,
"ret": round(cumret - 1, 6),
"np": round(net_p, 2),
"mdd": round(mdd, 6),
"pf": round(pf, 2),
"rr": round(rr, 2),
"expect": round(expect, 2),
}
def compute_metrics(bt_result, close_test, capital=10000):
"""Compute all standard metrics from backtest result.
Uses trade-level compounding (same as webapp _trade_stats) for accuracy.
Returns: dict with total_ret, bh_ret, sharpe_strat, sharpe_bh, mdd, n_trades
"""
equity = bt_result["equity"]
trade_returns = bt_result["trade_returns"]
# Total return — trade-level compounding (matches webapp)
if trade_returns:
cumret = 1.0
for r in trade_returns:
cumret *= (1 + r)
total_ret = cumret - 1
else:
total_ret = 0.0
# Buy and hold
bh_equity = capital * (close_test / close_test.iloc[0])
bh_ret = (bh_equity.iloc[-1] - capital) / capital if capital != 0 else 0.0
# Sharpe ratio — trade-level (matches webapp: sqrt(252*26) annualization)
if len(trade_returns) >= 2 and float(np.std(trade_returns)) > 0:
sharpe_strat = float(np.mean(trade_returns) / np.std(trade_returns) * np.sqrt(252 * 26))
else:
sharpe_strat = 0.0
bh_rets = bh_equity.pct_change().dropna()
if len(bh_rets) > 1 and bh_rets.std() != 0:
sharpe_bh = float((bh_rets.mean() / bh_rets.std()) * np.sqrt(252 * 24 * 4))
else:
sharpe_bh = 0.0
# Max drawdown — trade-level (matches webapp)
if trade_returns:
eq = np.cumprod([1.0] + [1 + r for r in trade_returns])
peak = np.maximum.accumulate(eq)
mdd = float(((eq - peak) / peak).min()) if len(eq) > 1 else 0.0
else:
mdd = 0.0
return {
"total_ret": float(total_ret),
"bh_ret": float(bh_ret),
"sharpe_strat": float(sharpe_strat) if not np.isnan(sharpe_strat) else 0.0,
"sharpe_bh": float(sharpe_bh) if not np.isnan(sharpe_bh) else 0.0,
"mdd": float(mdd),
"n_trades": len(trade_returns),
}
# Diagnostics line/histogram series (equity / drawdown / rolling_acc / conf_hist)
# only feed the small Diagnostics charts — they're never used by the price chart
# or scroll-back. On a 1-min model trained over the (2.2-capped) window these are
# still ~30k points each; downsample to a visually-identical resolution before the
# dict leaves the trainer so it doesn't carry that into Server-A RAM / Postgres.
_RESULTS_SERIES_MAX = 5000
def _downsample_idx(n, cap=_RESULTS_SERIES_MAX):
"""Evenly-spaced index list spanning [0, n-1] (first+last always kept), or
None when no downsampling is needed (n <= cap)."""
if n <= cap:
return None
return np.unique(np.linspace(0, n - 1, cap).astype(int)).tolist()
def _take(arr, idx):
"""Subset a list by an index list (idx may be None → return arr unchanged)."""
if idx is None or not isinstance(arr, list):
return arr
return [arr[i] for i in idx]
# trade_log / train_trade_log are lists of per-trade dicts (display-only — the
# Trade Log tab). They scale with TRADE count, not bar count, so the bar-window
# cap (Phase 2.2) doesn't bound them — a degenerate near-every-bar model can put
# 10k+ trade dicts in the blob (>3 MB). Cap each (independently — a small-N model
# keeps every trade) to the most-recent N, recording `*_total` + `*_truncated`
# so the true count is still reported. Real strategies have far fewer than
# _TRADE_LOG_MAX trades, so this only ever bites pathological models.
_TRADE_LOG_MAX = 5000
def _cap_trade_log(tl):
"""Return (capped_list, original_len, was_truncated)."""
if not isinstance(tl, list) or len(tl) <= _TRADE_LOG_MAX:
return tl, (len(tl) if isinstance(tl, list) else 0), False
return tl[-_TRADE_LOG_MAX:], len(tl), True
def build_return_dict(split_result, bt_result, metrics, model, feature_cols,
signal_full, p_pos_test, p_neg_test, custom_figs=None,
bt_train_result=None, pre_stats=None):
"""Assemble the complete return dict. Handles ALL serialization.
Never returns Timestamps, numpy arrays, or non-JSON types.
Returns: JSON-safe dict with all required keys
"""
df = split_result["df"]
close = split_result["close"]
close_test = split_result["close_test"]
X_test = split_result["X_test"]
y_test = split_result["y_test"]
equity = bt_result["equity"]
bar_returns = bt_result["bar_returns"]
# OHLC
ohlc_dates = [str(x) for x in df.index.tolist()]
def _safe_list(arr):
return [float(x) if (x is not None and not np.isnan(x) and not np.isinf(x)) else None
for x in arr]
# Overlays
bb, ma = compute_overlays(close, df.index)
# Buy and hold equity
capital = equity.iloc[0] if len(equity) > 0 else 10000
bh_equity = capital * (close_test / close_test.iloc[0])
# Confusion matrix
from sklearn.metrics import confusion_matrix
pred_test = model.predict(X_test)
y_test_arr = np.asarray(y_test)
cm = confusion_matrix(y_test_arr, pred_test, labels=[-1, 0, 1])
# Rolling accuracy
sig_arr = signal_full.reindex(close_test.index).values
correct = pd.Series((pred_test == y_test_arr).astype(float), index=X_test.index)
active_test = pd.Series(sig_arr != 0, index=close_test.index) if len(sig_arr) == len(close_test) else pd.Series(True, index=close_test.index)
correct_active = correct.where(active_test, other=np.nan)
rolling_acc = correct_active.rolling(30, min_periods=1).mean()
# Feature importance
importances = model.feature_importances_
fi_pairs = sorted(zip(feature_cols, importances), key=lambda x: x[1])[-15:]
# Drawdown
rolling_max = equity.cummax()
drawdown = (equity - rolling_max) / rolling_max.replace(0, np.nan)
drawdown = drawdown.fillna(0.0)
# ── Downsample the Diagnostics-only series (see _downsample_idx) ──────────
_eq_dates = [str(x) for x in close_test.index.tolist()]
_eq_strat = _safe_list(equity.values)
_eq_bh = _safe_list(bh_equity.values)
_eq_idx = _downsample_idx(len(_eq_dates))
_eq_dates, _eq_strat, _eq_bh = _take(_eq_dates, _eq_idx), _take(_eq_strat, _eq_idx), _take(_eq_bh, _eq_idx)
_ra_dates = [str(x) for x in rolling_acc.index.tolist()]
_ra_vals = [float(x) if (not np.isnan(x) and not np.isinf(x)) else None for x in rolling_acc.values]
_ra_idx = _downsample_idx(len(_ra_dates))
_ra_dates, _ra_vals = _take(_ra_dates, _ra_idx), _take(_ra_vals, _ra_idx)
_dd_dates = [str(x) for x in drawdown.index.tolist()]
_dd_vals = _safe_list(drawdown.values)
_dd_idx = _downsample_idx(len(_dd_dates))
_dd_dates, _dd_vals = _take(_dd_dates, _dd_idx), _take(_dd_vals, _dd_idx)
_cp_pos = [float(x) for x in (p_pos_test.tolist() if hasattr(p_pos_test, 'tolist') else list(p_pos_test))]
_cp_neg = [float(x) for x in (p_neg_test.tolist() if hasattr(p_neg_test, 'tolist') else list(p_neg_test))]
_cp_pos = _take(_cp_pos, _downsample_idx(len(_cp_pos)))
_cp_neg = _take(_cp_neg, _downsample_idx(len(_cp_neg)))
# ── Trade logs — display-only (Trade Log tab); cap to most-recent N with a
# `_total` field so the true count is still reported (see _cap_trade_log).
# NB: ret_dist arrays are left FULL — a downstream path in callbacks.py
# recomputes n_trades/win-rate from len(ret_dist), so a sample would skew
# the displayed counts; they're small anyway and gzip handles them.
_tl_test, _tl_test_n, _tl_test_tr = _cap_trade_log(bt_result.get("trade_log", []))
_tl_tr, _tl_tr_n, _tl_tr_tr = _cap_trade_log(bt_train_result.get("trade_log", []) if bt_train_result else [])
return {
"ohlc": {
"dates": ohlc_dates,
"open": _safe_list(df["open"].values),
"high": _safe_list(df["high"].values),
"low": _safe_list(df["low"].values),
"close": _safe_list(df["close"].values),
},
"signals": {
"dates": [str(x) for x in signal_full.index.tolist()],
"values": [float(x) for x in signal_full.values],
},
"bb": bb,
"ma": ma,
"equity": {
"dates": _eq_dates,
"strategy": _eq_strat,
"bh": _eq_bh,
},
"feature_importance": {
"names": [p[0] for p in fi_pairs],
"values": [float(p[1]) for p in fi_pairs],
},
"conf_matrix": cm.tolist(),
"conf_hist": {
"p_pos": _cp_pos,
"p_neg": _cp_neg,
},
"rolling_acc": {
"dates": _ra_dates,
"values": _ra_vals,
},
"drawdown": {
"dates": _dd_dates,
"values": _dd_vals,
},
"ret_dist": [float(x) for x in bt_result["trade_returns"]],
"ret_dist_long": [float(x) for x in bt_result["long_returns"]],
"ret_dist_short": [float(x) for x in bt_result["short_returns"]],
"train_ret_dist": [float(x) for x in bt_train_result["trade_returns"]] if bt_train_result else [],
"train_ret_dist_long": [float(x) for x in bt_train_result["long_returns"]] if bt_train_result else [],
"train_ret_dist_short": [float(x) for x in bt_train_result["short_returns"]] if bt_train_result else [],
"trade_log": _tl_test,
"train_trade_log": _tl_tr,
"trade_log_total": _tl_test_n,
"train_trade_log_total": _tl_tr_n,
"trade_log_truncated": _tl_test_tr,
"train_trade_log_truncated": _tl_tr_tr,
**(pre_stats or {}),
"metrics": metrics,
"split_dt": split_result["split_dt"],
"split_idx": int(split_result["split_idx"]),
"n_train": int(split_result["n_train"]),
"n_test": int(split_result["n_test"]),
"feature_cols": list(feature_cols),
"custom_figs": custom_figs or [],
}
# ════════════════════════════════════════════════════════════════════════════
# STRATEGY FRAMEWORK v2 — Config-driven architecture
# Claude writes feature_engineering() + strategy_config(). Framework does rest.
# ════════════════════════════════════════════════════════════════════════════
import importlib
_MODEL_REGISTRY = {
"XGBClassifier": ("xgboost", "XGBClassifier"),
"RandomForestClassifier": ("sklearn.ensemble", "RandomForestClassifier"),
"GradientBoostingClassifier": ("sklearn.ensemble", "GradientBoostingClassifier"),
"LogisticRegression": ("sklearn.linear_model", "LogisticRegression"),
"ExtraTreesClassifier": ("sklearn.ensemble", "ExtraTreesClassifier"),
"AdaBoostClassifier": ("sklearn.ensemble", "AdaBoostClassifier"),
}
def _build_model_from_config(config, X_train, y_train_enc):
"""Build, fit, and wrap a model from strategy_config dict."""
model_type = config.get("model_type", "RandomForestClassifier")
model_params = dict(config.get("model_params", {}))
if model_type not in _MODEL_REGISTRY:
raise ValueError(f"Unknown model_type '{model_type}'. Valid: {list(_MODEL_REGISTRY.keys())}")
module_path, class_name = _MODEL_REGISTRY[model_type]
mod = importlib.import_module(module_path)
cls = getattr(mod, class_name)
# XGBoost defaults
if class_name == "XGBClassifier":
model_params.setdefault("use_label_encoder", False)
model_params.setdefault("eval_metric", "mlogloss")
model_params.setdefault("tree_method", "hist")
# Determinism > speed (2026-05-25). XGBoost hist with n_jobs=-1 is
# NON-reproducible even with random_state set — the parallel histogram
# gradient-sum order varies across threads, so the SAME code + data
# gives a slightly different model (and backtest) every run. Forcing
# single-thread makes training bit-reproducible so: (a) a user who
# copies a strategy and reruns it gets identical numbers, (b) the
# community "Live" score matches a redeploy, (c) "same code, different
# result" support reports go away. Cost: single-threaded XGB (a few
# seconds slower on large windows; hist is fast so it's minor). FORCED
# (not setdefault) so the guarantee can't be silently broken by a
# strategy passing n_jobs. Exact reproducibility holds within the
# platform (pinned versions / same Modal image); a user's own machine
# with different xgboost/numpy/CPU can still differ in low-order bits.
model_params["n_jobs"] = 1
# Common defaults
model_params.setdefault("random_state", 42)
from model_wrapper import ModelWrapper
clf = cls(**model_params)
clf.fit(X_train, y_train_enc)
enc = LabelEncoder()
enc.fit([-1, 0, 1])
return ModelWrapper(clf, original_classes=enc.classes_, n_features=X_train.shape[1])
def _generate_signals(model, X, threshold):
"""Framework-owned signal generation. Deterministic threshold logic."""
proba = model.predict_proba(X)
classes = list(model.classes_)
idx_pos = classes.index(1) if 1 in classes else None
idx_neg = classes.index(-1) if -1 in classes else None
p_pos = proba[:, idx_pos] if idx_pos is not None else np.zeros(len(X))
p_neg = proba[:, idx_neg] if idx_neg is not None else np.zeros(len(X))
signal_vals = np.zeros(len(X))
signal_vals = np.where(p_pos >= threshold, 1.0, signal_vals)
signal_vals = np.where(p_neg >= threshold, -1.0, signal_vals)
# Both exceed: pick stronger
both = (p_pos >= threshold) & (p_neg >= threshold)
signal_vals[both] = np.where(p_pos[both] >= p_neg[both], 1.0, -1.0)
return pd.Series(signal_vals, index=X.index), p_pos, p_neg
# ── Filter functions (all no-ops when config value is None) ──────────────
def _apply_direction_filter(signal, direction):
"""Zero out signals that don't match allowed direction."""
if direction is None or direction == "both":
return signal
s = signal.copy()
if direction == "long":
s[s < 0] = 0.0
elif direction == "short":
s[s > 0] = 0.0
return s
def _apply_session_filter(signal, index, session_hours):
"""Zero out signals outside session hours [start, end] UTC."""
if session_hours is None:
return signal
s = signal.copy()
start_h, end_h = session_hours[0], session_hours[1]
hours = index.hour
if start_h <= end_h:
mask = (hours >= start_h) & (hours < end_h)
else: # wrap around midnight, e.g. [22, 6]
mask = (hours >= start_h) | (hours < end_h)
s[~mask] = 0.0
return s
def _apply_atr_filter(signal, close, high, low, min_atr):
"""Zero out signals when NATR(14) is below threshold."""
if min_atr is None:
return signal
hl = high - low
hc = (high - close.shift(1)).abs()
lc = (low - close.shift(1)).abs()
tr = pd.concat([hl, hc, lc], axis=1).max(axis=1)
atr14 = tr.ewm(com=13, adjust=False).mean()
natr = atr14 / close.replace(0, np.nan)
s = signal.copy()
s[natr < min_atr] = 0.0
return s
def _apply_trend_filter(signal, close, trend_filter):
"""Only allow signals aligned with trend. e.g. 'sma_50': longs above SMA, shorts below."""
if trend_filter is None:
return signal
# Parse: "sma_50" → SMA with period 50
parts = trend_filter.lower().replace("-", "_").split("_")
if len(parts) >= 2 and parts[0] in ("sma", "ema"):
period = int(parts[1])
else:
return signal # unknown filter, skip
if parts[0] == "sma":
trend_line = close.rolling(period).mean()
else:
trend_line = close.ewm(span=period, adjust=False).mean()
s = signal.copy()
# Longs only above trend, shorts only below
s[(s > 0) & (close < trend_line)] = 0.0
s[(s < 0) & (close > trend_line)] = 0.0
return s
# ── run_backtest_v2: framework-owned SL/TP/cooldown/position management ──
def run_backtest_v2(signal, close, high, low, config, capital=10000, cost=2e-5):
"""Backtest with SL/TP/cooldown/direction handling built into the engine.
Unlike run_backtest (v1), this function handles position exits internally.
Returns: same dict shape as run_backtest()
"""
stop_loss = config.get("stop_loss")
take_profit = config.get("take_profit")
cooldown = config.get("cooldown", 0)
on_opposite = config.get("on_opposite", "reverse")
sig_arr = signal.values
close_arr = close.values
high_arr = high.values
low_arr = low.values
idx = signal.index
n = len(close_arr)
trade_returns = []
long_returns = []
short_returns = []
trade_log = []
equity_vals = np.full(n, float(capital))
cumret = 1.0
position = 0.0 # current direction: 1.0, -1.0, or 0.0 (flat)
entry_price = None
entry_bar = None # index into arrays for entry time
cooldown_remaining = 0
def _log_trade(exit_bar, exit_px, ret, reason):
trade_log.append({
"type": "Buy" if position == 1.0 else "Sell",
"entry_time": str(idx[entry_bar]),
"exit_time": str(idx[exit_bar]),
"entry_price": round(entry_price, 5),
"exit_price": round(exit_px, 5),
"pnl": round(position * (exit_px - entry_price), 5),
"pnl_pct": round(ret * 100, 3),
"exit_reason": reason,
})
for i in range(n):
c = close_arr[i]
h = high_arr[i]
lo = low_arr[i]
s = sig_arr[i]
# 1. Check SL/TP if in trade
if position != 0.0 and entry_price is not None:
hit_sl = False
hit_tp = False
exit_price = None
if position == 1.0: # long
if stop_loss is not None and lo <= entry_price * (1 - stop_loss):
hit_sl = True
exit_price = entry_price * (1 - stop_loss)
elif take_profit is not None and h >= entry_price * (1 + take_profit):
hit_tp = True
exit_price = entry_price * (1 + take_profit)
else: # short
if stop_loss is not None and h >= entry_price * (1 + stop_loss):
hit_sl = True
exit_price = entry_price * (1 + stop_loss)
elif take_profit is not None and lo <= entry_price * (1 - take_profit):
hit_tp = True
exit_price = entry_price * (1 - take_profit)
if hit_sl or hit_tp:
ret = float(position * (exit_price - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, exit_price, ret, "SL" if hit_sl else "TP")
cumret *= (1 + ret)
position = 0.0
entry_price = None
entry_bar = None
cooldown_remaining = cooldown
equity_vals[i] = capital * cumret
continue
# 2. Cooldown
if cooldown_remaining > 0:
cooldown_remaining -= 1
equity_vals[i] = capital * cumret
continue
# 3. Signal processing
if s != 0.0:
if position == 0.0:
# Open new trade
position = s
entry_price = c
entry_bar = i
elif s != position:
# Opposite signal
if on_opposite == "reverse":
# Close current + open opposite
ret = float(position * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, c, ret, "signal")
cumret *= (1 + ret)
position = s
entry_price = c
entry_bar = i
else: # close_only
# Close current, go flat
ret = float(position * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(i, c, ret, "close_only")
cumret *= (1 + ret)
position = 0.0
entry_price = None
entry_bar = None
cooldown_remaining = cooldown
equity_vals[i] = capital * cumret
# Close last open trade at final close
if position != 0.0 and entry_price is not None and n > 0 and entry_price != 0:
c = close_arr[-1]
ret = float(position * (c - entry_price) / entry_price - cost)
trade_returns.append(ret)
if position == 1.0:
long_returns.append(ret)
else:
short_returns.append(ret)
_log_trade(n - 1, c, ret, "end")
cumret *= (1 + ret)
equity_vals[-1] = capital * cumret
# Bar returns for Sharpe (approximate)
bar_returns = np.zeros(n)
for i in range(1, n):
if close_arr[i - 1] != 0 and sig_arr[i - 1] != 0:
bar_returns[i] = sig_arr[i - 1] * (close_arr[i] - close_arr[i - 1]) / close_arr[i - 1]
return {
"equity": pd.Series(equity_vals, index=close.index),
"trade_returns": trade_returns,
"long_returns": long_returns,
"short_returns": short_returns,
"bar_returns": bar_returns,
"trade_log": trade_log,
}
# ── run_strategy: the v2 orchestrator ────────────────────────────────────
def run_strategy(feature_fn, config_fn, data_path, start_date="", end_date="",
validation_date="", train_split=0.7, register_model_fn=None):
"""Config-driven strategy execution. Claude writes feature_fn + config_fn,
framework does everything else.
Returns: results dict (same format as webapp expects)
"""
config = config_fn()
# Auto-correct SL/TP if Claude passed percentage instead of decimal
for _key in ("stop_loss", "take_profit"):
_val = config.get(_key)
if _val is not None and _val > 0.1: # >10% is almost certainly a percentage
config[_key] = _val / 100.0
print(f"[strategy] Auto-corrected {_key}: {_val} -> {config[_key]} (was percentage, converted to decimal)")
# 1. Load data
df, close, open_, high, low = load_ohlc(data_path, start_date, end_date)
# 2. Feature engineering (Claude's function)
df = feature_fn(df, close, open_, high, low)
close = df["close"]
open_ = df["open"]
high = df["high"]
low = df["low"]
# 3. Warm-up detection: drop rows where features have NaN BEFORE any fill
feature_cols = [c for c in df.columns if c not in ("open", "high", "low", "close")]
raw_nans = df[feature_cols].isna().any(axis=1)
valid_rows = ~raw_nans
if valid_rows.any():
first_valid = valid_rows.idxmax()
if raw_nans.loc[:first_valid].any():
df = df.loc[first_valid:].copy()
close = df["close"]
open_ = df["open"]
high = df["high"]
low = df["low"]
# 4. Target
horizon = config.get("target_horizon", 4)
target = make_target(close, horizon=horizon)
# 5. Split (ffill only within each partition — no bfill leak)
mask = target.notna()
df = df[mask].copy()
target = target[mask]
close = df["close"]
high = df["high"]
low = df["low"]
X = df[feature_cols].copy()
X = X.replace([np.inf, -np.inf], np.nan)
if validation_date:
split_idx = len(df[df.index <= validation_date])
else:
split_idx = int(len(df) * train_split)
split_idx = max(1, min(split_idx, len(df) - 1))
# ffill within train and test separately (no leak)
X_train = X.iloc[:split_idx].ffill().fillna(0.0)
X_test = X.iloc[split_idx:].ffill().fillna(0.0)
X = pd.concat([X_train, X_test])
y_train = target.iloc[:split_idx]
y_test = target.iloc[split_idx:]
close_train = close.iloc[:split_idx]
close_test = close.iloc[split_idx:]
high_test = high.iloc[split_idx:]
low_test = low.iloc[split_idx:]
enc = LabelEncoder()
enc.fit([-1, 0, 1])
y_train_enc = enc.transform(y_train)
y_test_enc = enc.transform(y_test)
split_dt = str(df.index[split_idx])
sp = {
"df": df, "X_train": X_train, "X_test": X_test,
"y_train": y_train, "y_test": y_test,
"y_train_enc": y_train_enc, "y_test_enc": y_test_enc,
"enc": enc,
"close": close, "close_train": close_train, "close_test": close_test,
"split_idx": split_idx, "split_dt": split_dt,
"n_train": len(X_train), "n_test": len(X_test),
}
# 6. Build model from config
model = _build_model_from_config(config, X_train, y_train_enc)
# 7. Generate signals
threshold = config.get("signal_threshold", 0.55)
signal_train, p_pos_train, p_neg_train = _generate_signals(model, X_train, threshold)
signal_test, p_pos_test, p_neg_test = _generate_signals(model, X_test, threshold)
# 8. Apply filters (order: direction → session → ATR → trend)
direction = config.get("direction", "both")
signal_test = _apply_direction_filter(signal_test, direction)
signal_train = _apply_direction_filter(signal_train, direction)
session_filter = config.get("session_filter")
signal_test = _apply_session_filter(signal_test, signal_test.index, session_filter)
signal_train = _apply_session_filter(signal_train, signal_train.index, session_filter)
min_atr = config.get("min_atr")
if min_atr is not None:
signal_test = _apply_atr_filter(signal_test, close_test, high_test, low_test, min_atr)
trend_filter = config.get("trend_filter")
if trend_filter is not None:
signal_test = _apply_trend_filter(signal_test, close_test, trend_filter)
signal_full = pd.concat([signal_train, signal_test])
# 9. Backtest with SL/TP/cooldown (test + train)
high_train = high.iloc[:split_idx]
low_train = low.iloc[:split_idx]
has_risk = (config.get("stop_loss") is not None or
config.get("take_profit") is not None or
config.get("cooldown", 0) > 0 or
config.get("on_opposite", "reverse") != "reverse")
if has_risk:
bt = run_backtest_v2(signal_test, close_test, high_test, low_test, config, capital=10000)
bt_train = run_backtest_v2(signal_train, close_train, high_train, low_train, config, capital=10000)
else:
bt = run_backtest(signal_test, close_test, capital=10000)
bt_train = run_backtest(signal_train, close_train, capital=10000)
# 10. Metrics
metrics = compute_metrics(bt, close_test, capital=10000)
# 11. Pre-compute all trade stats (single source of truth)
pre_stats = {
"train_stats": compute_trade_stats(bt_train.get("trade_returns", []), capital=10000),
"test_stats": compute_trade_stats(bt.get("trade_returns", []), capital=10000),
"long_stats": compute_trade_stats(bt.get("long_returns", []), capital=10000),
"short_stats": compute_trade_stats(bt.get("short_returns", []), capital=10000),
}
# 12. Register model
if register_model_fn is not None:
register_model_fn(model)
# 13. Build return dict
return build_return_dict(sp, bt, metrics, model, feature_cols,
signal_full, p_pos_test, p_neg_test, custom_figs=[],
bt_train_result=bt_train, pre_stats=pre_stats)
# ── End strategy_utils ──
DATA_PATH = '/root/Desktop/QuantifyMe/data/ohlc/EURUSD_15min.parquet'
START_DATE = '2026-04-15'
END_DATE = '2026-05-25'
VALIDATION_DATE = ""
TRAIN_SPLIT = 0.7
# SECTION 1 — FEATURE ENGINEERING
def feature_engineering(df, close, open_, high, low):
# --- RSI(14) ---
period = 14
delta = close.diff()
gain = delta.clip(lower=0)
loss = -delta.clip(upper=0)
avg_gain = gain.ewm(alpha=1.0 / period, min_periods=period, adjust=False).mean()
avg_loss = loss.ewm(alpha=1.0 / period, min_periods=period, adjust=False).mean()
rs = avg_gain / avg_loss.replace(0, np.nan)
df["rsi_14"] = 100.0 - (100.0 / (1.0 + rs))
# --- RSI crossover signals: cross up from below 30, cross down from above 70 ---
rsi_prev = df["rsi_14"].shift(1)
df["rsi_cross_up30"] = np.where(
(rsi_prev < 30) & (df["rsi_14"] >= 30), 1.0, 0.0
)
df["rsi_cross_dn70"] = np.where(
(rsi_prev > 70) & (df["rsi_14"] <= 70), 1.0, 0.0
)
# --- RSI distance from thresholds (signed) ---
df["rsi_dist_30"] = df["rsi_14"] - 30.0
df["rsi_dist_70"] = df["rsi_14"] - 70.0
df["rsi_dist_50"] = df["rsi_14"] - 50.0
# --- RSI(5) for short-term momentum ---
period5 = 5
delta5 = close.diff()
gain5 = delta5.clip(lower=0)
loss5 = -delta5.clip(upper=0)
avg_gain5 = gain5.ewm(alpha=1.0 / period5, min_periods=period5, adjust=False).mean()
avg_loss5 = loss5.ewm(alpha=1.0 / period5, min_periods=period5, adjust=False).mean()
rs5 = avg_gain5 / avg_loss5.replace(0, np.nan)
df["rsi_5"] = 100.0 - (100.0 / (1.0 + rs5))
# --- RSI(28) for longer-term regime ---
period28 = 28
delta28 = close.diff()
gain28 = delta28.clip(lower=0)
loss28 = -delta28.clip(upper=0)
avg_gain28 = gain28.ewm(alpha=1.0 / period28, min_periods=period28, adjust=False).mean()
avg_loss28 = loss28.ewm(alpha=1.0 / period28, min_periods=period28, adjust=False).mean()
rs28 = avg_gain28 / avg_loss28.replace(0, np.nan)
df["rsi_28"] = 100.0 - (100.0 / (1.0 + rs28))
# --- Bollinger Bands (20, 2) ---
bb_period = 20
bb_mid = close.rolling(bb_period).mean()
bb_std = close.rolling(bb_period).std()
bb_upper = bb_mid + 2.0 * bb_std
bb_lower = bb_mid - 2.0 * bb_std
df["bb_mid"] = bb_mid
df["bb_width"] = np.where(bb_mid != 0, (bb_upper - bb_lower) / bb_mid, np.nan)
df["bb_pct_b"] = np.where(
(bb_upper - bb_lower) != 0,
(close - bb_lower) / (bb_upper - bb_lower),
0.5
)
df["price_vs_bb_mid"] = close - bb_mid
# --- ATR(14) for volatility ---
tr1 = high - low
tr2 = (high - close.shift(1)).abs()
tr3 = (low - close.shift(1)).abs()
true_range = pd.concat([tr1, tr2, tr3], axis=1).max(axis=1)
df["atr_14"] = true_range.ewm(alpha=1.0 / 14, min_periods=14, adjust=False).mean()
df["natr_14"] = np.where(close != 0, df["atr_14"] / close, np.nan)
# --- MACD (12, 26, 9) ---
ema12 = close.ewm(span=12, min_periods=12, adjust=False).mean()
ema26 = close.ewm(span=26, min_periods=26, adjust=False).mean()
macd_line = ema12 - ema26
signal_line = macd_line.ewm(span=9, min_periods=9, adjust=False).mean()
df["macd"] = macd_line
df["macd_signal"] = signal_line
df["macd_hist"] = macd_line - signal_line
df["macd_cross_up"] = np.where(
(macd_line.shift(1) < signal_line.shift(1)) & (macd_line >= signal_line), 1.0, 0.0
)
df["macd_cross_dn"] = np.where(
(macd_line.shift(1) > signal_line.shift(1)) & (macd_line <= signal_line), 1.0, 0.0
)
# --- Stochastic Oscillator (14, 3) ---
stoch_period = 14
lowest_low = low.rolling(stoch_period).min()
highest_high = high.rolling(stoch_period).max()
denom = (highest_high - lowest_low).replace(0, np.nan)
stoch_k = 100.0 * (close - lowest_low) / denom
stoch_d = stoch_k.rolling(3).mean()
df["stoch_k"] = stoch_k
df["stoch_d"] = stoch_d
df["stoch_cross_up"] = np.where(
(stoch_k.shift(1) < stoch_d.shift(1)) & (stoch_k >= stoch_d) & (stoch_k < 30), 1.0, 0.0
)
df["stoch_cross_dn"] = np.where(
(stoch_k.shift(1) > stoch_d.shift(1)) & (stoch_k <= stoch_d) & (stoch_k > 70), 1.0, 0.0
)
# --- SMA trend features ---
df["sma_20"] = close.rolling(20).mean()
df["sma_50"] = close.rolling(50).mean()
df["sma_200"] = close.rolling(200).mean()
df["price_vs_sma20"] = (close - df["sma_20"]) / df["sma_20"].replace(0, np.nan)
df["price_vs_sma50"] = (close - df["sma_50"]) / df["sma_50"].replace(0, np.nan)
df["sma20_vs_sma50"] = (df["sma_20"] - df["sma_50"]) / df["sma_50"].replace(0, np.nan)
# --- Price momentum (returns) ---
df["ret_1"] = close.pct_change(1)
df["ret_4"] = close.pct_change(4)
df["ret_8"] = close.pct_change(8)
df["ret_16"] = close.pct_change(16)
# --- Candle body and wick features ---
body = (close - open_).abs()
upper_wick = high - pd.concat([close, open_], axis=1).max(axis=1)
lower_wick = pd.concat([close, open_], axis=1).min(axis=1) - low
candle_range = (high - low).replace(0, np.nan)
df["body_ratio"] = body / candle_range
df["upper_wick_ratio"] = upper_wick / candle_range
df["lower_wick_ratio"] = lower_wick / candle_range
df["candle_direction"] = np.where(close >= open_, 1.0, -1.0)
# --- Volume of consecutive bars in same direction ---
df["consec_up"] = (
df["candle_direction"]
.groupby((df["candle_direction"] != df["candle_direction"].shift(1)).cumsum())
.cumcount() + 1
) * np.where(df["candle_direction"] > 0, 1.0, 0.0)
df["consec_dn"] = (
df["candle_direction"]
.groupby((df["candle_direction"] != df["candle_direction"].shift(1)).cumsum())
.cumcount() + 1
) * np.where(df["candle_direction"] < 0, 1.0, 0.0)
# --- RSI slope (rate of change) ---
df["rsi_slope_3"] = df["rsi_14"].diff(3)
df["rsi_slope_5"] = df["rsi_14"].diff(5)
# --- Rolling high/low channel breakout context ---
df["high_20"] = high.rolling(20).max()
df["low_20"] = low.rolling(20).min()
df["price_pos_in_range"] = np.where(
(df["high_20"] - df["low_20"]) != 0,
(close - df["low_20"]) / (df["high_20"] - df["low_20"]),
0.5
)
# --- RSI oversold/overbought binary flags ---
df["rsi_oversold"] = np.where(df["rsi_14"] < 30, 1.0, 0.0)
df["rsi_overbought"] = np.where(df["rsi_14"] > 70, 1.0, 0.0)
df["rsi_neutral"] = np.where((df["rsi_14"] >= 40) & (df["rsi_14"] <= 60), 1.0, 0.0)
# --- Rolling RSI min/max to track extremes ---
df["rsi_min_10"] = df["rsi_14"].rolling(10).min()
df["rsi_max_10"] = df["rsi_14"].rolling(10).max()
df["rsi_range_10"] = df["rsi_max_10"] - df["rsi_min_10"]
# --- Time-of-day features (sin/cos encoding for session awareness) ---
if hasattr(df.index, "hour"):
hour = df.index.hour + df.index.minute / 60.0
df["hour_sin"] = np.sin(2.0 * np.pi * hour / 24.0)
df["hour_cos"] = np.cos(2.0 * np.pi * hour / 24.0)
else:
df["hour_sin"] = 0.0
df["hour_cos"] = 1.0
# --- Fill NaN from warm-up periods ---
df = df.bfill().ffill()
return df
# SECTION 2 — STRATEGY CONFIG
def strategy_config():
return {
"title": "RSI Crossover Mean-Reversion (XGBoost, Sharpe)",
"model_type": "XGBClassifier",
"model_params": {
"n_estimators": 400,
"max_depth": 4,
"learning_rate": 0.04,
"subsample": 0.8,
"colsample_bytree": 0.7,
"min_child_weight": 3,
"gamma": 0.1,
"reg_alpha": 0.05,
"reg_lambda": 1.2,
"objective": "binary:logistic",
"n_jobs": -1,
"random_state": 42,
},
"signal_threshold": 0.55,
"direction": "both",
"stop_loss": 0.004,
"take_profit": 0.008,
"cooldown": 0,
"max_positions": 1,
"on_opposite": "reverse",
"session_filter": [7, 18],
"min_atr": None,
"trend_filter": None,
"target_horizon": 4,
"objective": (
"Maximize Sharpe ratio by capturing mean-reversion when RSI crosses "
"back from oversold (<30) or overbought (>70) extremes. XGBoost with "
"moderate depth and shrinkage prevents overfitting on the short EUR/USD "
"window. A 2:1 TP:SL ratio (0.8%/0.4%) on 15-min bars targets clean "
"risk-adjusted returns. Session filter restricts to liquid London/NY hours."
),
"notes": (
"Feature set combines the core RSI crossover signal with multi-period RSI, "
"MACD histogram, Stochastic, Bollinger %B, ATR normalised volatility, "
"price momentum across 4 horizons, candle structure ratios, and time encoding. "
"colsample_bytree=0.7 adds diversity across trees; subsample=0.8 reduces "
"variance. min_child_weight=3 avoids splitting on noisy one-off RSI spikes. "
"No trend_filter so the model can express both long and short mean-reversion "
"signals symmetrically via the 'both' direction setting."
),
}
# ── Framework v2: auto-generated wrapper ──
def train_and_backtest():
_vd = VALIDATION_DATE if 'VALIDATION_DATE' in globals() else ''
_ts = TRAIN_SPLIT if 'TRAIN_SPLIT' in globals() else 0.7
return run_strategy(
feature_engineering, strategy_config,
DATA_PATH, START_DATE, END_DATE,
_vd, _ts,
register_model_fn=register_model
)
|
||||||||||
|
—
|
USD/JPY Multi-MA + RSI/BB XGBoost Sharpe
Maximize Sharpe ratio on USD/JPY 1-min data using XGBoost with returns, RSI, Bollinger Bands, multiple MAs (50/100/200), MACD, ATR, and cand…
|
M
@malcolmtan
|
USDJPY | 1min | 60.7%44.3% | +0.42%-4.25% | 1.220.84 | 0.53%0.53% | 84106 |
|
# ╔══════════════════════════════════════════════════════════════╗
# ║ STRATEGY REQUEST LOG ║
# ╚══════════════════════════════════════════════════════════════╝
# Generated : 2026-05-08 02:08:02
# Model : XGBoost
# Feature Eng. : Auto-add features: ON
# Signal / Entry : —
# Optimization : —
# Risk Mgmt : —
# Risk Filter : —
# ══════════════════════════════════════════════════════════════
# ============================================================
# SECTION 0 — IMPORTS & CONSTANTS
import numpy as np
import pandas as pd
DATA_PATH = "/root/Desktop/QuantifyMe/data/ohlc/USDJPY_1min.parquet"
START_DATE = "2026-05-04 00:00:00"
END_DATE = "2026-05-07 00:00:00"
VALIDATION_DATE = ""
TRAIN_SPLIT = 0.6993736951983298
# SECTION 1 — FEATURE ENGINEERING
def feature_engineering(df, close, open_, high, low):
# --- Returns over multiple horizons ---
for n in [1, 3, 5, 10, 20]:
df[f"ret_{n}"] = close.pct_change(n)
# --- RSI 14 ---
delta = close.diff()
gain = delta.clip(lower=0)
loss = -delta.clip(upper=0)
avg_gain = gain.ewm(com=13, min_periods=14).mean()
avg_loss = loss.ewm(com=13, min_periods=14).mean()
rs = avg_gain / (avg_loss + 1e-10)
df["rsi_14"] = 100.0 - (100.0 / (1.0 + rs))
# --- RSI derived features ---
df["rsi_14_zscore"] = (df["rsi_14"] - df["rsi_14"].rolling(50).mean()) / (df["rsi_14"].rolling(50).std() + 1e-10)
df["rsi_overbought"] = np.where(df["rsi_14"] > 70, 1, 0)
df["rsi_oversold"] = np.where(df["rsi_14"] < 30, 1, 0)
# --- Bollinger Bands 20, 2 ---
bb_mid = close.rolling(20).mean()
bb_std = close.rolling(20).std()
bb_upper = bb_mid + 2.0 * bb_std
bb_lower = bb_mid - 2.0 * bb_std
df["bb_mid"] = bb_mid
df["bb_upper"] = bb_upper
df["bb_lower"] = bb_lower
df["bb_width"] = (bb_upper - bb_lower) / (bb_mid + 1e-10)
df["bb_pct_b"] = (close - bb_lower) / (bb_upper - bb_lower + 1e-10)
df["bb_above"] = np.where(close > bb_upper, 1, 0)
df["bb_below"] = np.where(close < bb_lower, 1, 0)
# --- Moving Averages ---
for w in [50, 100, 200]:
df[f"sma_{w}"] = close.rolling(w).mean()
df[f"price_vs_sma_{w}"] = (close - df[f"sma_{w}"]) / (df[f"sma_{w}"] + 1e-10)
# --- MA crossover signals ---
df["sma50_vs_sma100"] = np.where(df["sma_50"] > df["sma_100"], 1, -1)
df["sma50_vs_sma200"] = np.where(df["sma_50"] > df["sma_200"], 1, -1)
df["sma100_vs_sma200"] = np.where(df["sma_100"] > df["sma_200"], 1, -1)
# --- ATR 14 ---
tr = pd.concat([
high - low,
(high - close.shift(1)).abs(),
(low - close.shift(1)).abs()
], axis=1).max(axis=1)
df["atr_14"] = tr.ewm(com=13, min_periods=14).mean()
df["natr_14"] = df["atr_14"] / (close + 1e-10)
# --- Momentum / rate of change ---
for n in [5, 10, 20]:
df[f"mom_{n}"] = close - close.shift(n)
df[f"roc_{n}"] = (close - close.shift(n)) / (close.shift(n) + 1e-10)
# --- Volume features (if volume exists) ---
if "volume" in df.columns:
vol = df["volume"].replace(0, np.nan)
df["vol_sma_20"] = vol.rolling(20).mean()
df["vol_ratio_20"] = vol / (df["vol_sma_20"] + 1e-10)
else:
df["vol_ratio_20"] = 1.0
# --- Price spread & body features ---
df["hl_spread"] = (high - low) / (close + 1e-10)
df["body_ratio"] = (close - open_).abs() / (high - low + 1e-10)
df["upper_wick"] = (high - pd.concat([close, open_], axis=1).max(axis=1)) / (high - low + 1e-10)
df["lower_wick"] = (pd.concat([close, open_], axis=1).min(axis=1) - low) / (high - low + 1e-10)
df["bull_candle"] = np.where(close > open_, 1, 0)
# --- Lagged returns for autocorrelation signal ---
for lag in [1, 2, 3, 5]:
df[f"ret1_lag{lag}"] = df["ret_1"].shift(lag)
# --- Rolling volatility ---
df["vol_10"] = df["ret_1"].rolling(10).std()
df["vol_20"] = df["ret_1"].rolling(20).std()
df["vol_50"] = df["ret_1"].rolling(50).std()
# --- Z-score of close over 20 and 50 bars ---
df["zscore_20"] = (close - close.rolling(20).mean()) / (close.rolling(20).std() + 1e-10)
df["zscore_50"] = (close - close.rolling(50).mean()) / (close.rolling(50).std() + 1e-10)
# --- Relative distance of price from BB bands ---
df["dist_upper"] = (bb_upper - close) / (close + 1e-10)
df["dist_lower"] = (close - bb_lower) / (close + 1e-10)
# --- EMA 9 and 21 for short-term momentum ---
df["ema_9"] = close.ewm(span=9, min_periods=9).mean()
df["ema_21"] = close.ewm(span=21, min_periods=21).mean()
df["ema9_vs_ema21"] = np.where(df["ema_9"] > df["ema_21"], 1, -1)
df["price_vs_ema9"] = (close - df["ema_9"]) / (df["ema_9"] + 1e-10)
df["price_vs_ema21"] = (close - df["ema_21"]) / (df["ema_21"] + 1e-10)
# --- MACD-like signal ---
ema_12 = close.ewm(span=12, min_periods=12).mean()
ema_26 = close.ewm(span=26, min_periods=26).mean()
macd_line = ema_12 - ema_26
signal_line = macd_line.ewm(span=9, min_periods=9).mean()
df["macd"] = macd_line
df["macd_signal"] = signal_line
df["macd_hist"] = macd_line - signal_line
df["macd_cross"] = np.where(macd_line > signal_line, 1, -1)
# --- Fill NaN from warm-up periods ---
df = df.bfill().ffill()
return df
# SECTION 2 — STRATEGY CONFIG
def strategy_config():
return {
"title": "USD/JPY Multi-MA + RSI/BB XGBoost Sharpe",
"model_type": "XGBClassifier",
"model_params": {
"n_estimators": 400,
"max_depth": 4,
"learning_rate": 0.04,
"subsample": 0.75,
"colsample_bytree": 0.75,
"min_child_weight": 5,
"gamma": 0.1,
"reg_alpha": 0.1,
"reg_lambda": 1.5,
"objective": "binary:logistic",
"tree_method": "hist",
"random_state": 42,
},
"signal_threshold": 0.55,
"direction": "both",
"stop_loss": 0.0008,
"take_profit": 0.0016,
"cooldown": 0,
"max_positions": 1,
"on_opposite": "close_only",
"session_filter": None,
"min_atr": None,
"trend_filter": None,
"target_horizon": 5,
"objective": (
"Maximize Sharpe ratio on USD/JPY 1-min data using XGBoost with "
"returns, RSI, Bollinger Bands, multiple MAs (50/100/200), MACD, "
"ATR, and candle-body features. Stop-loss and take-profit set at "
"a 1:2 risk/reward to filter noise and improve Sharpe. n_estimators "
"and moderate depth balance bias-variance. Regularization (alpha/lambda) "
"reduces overfitting on short date range."
),
"notes": (
"Target horizon of 5 bars (5 minutes) is chosen to capture short-term "
"directional moves on 1-min data without excessive label noise. "
"colsample_bytree and subsample add stochasticity to reduce variance. "
"close_only on opposite signal avoids whipsaw from rapid reversals. "
"No session filter applied since USD/JPY has liquidity around the clock."
),
}
|
||||||||||