| """ |
| Walk-forward, look-ahead-safe backtest engine (spec sections 28-30, |
| 53-54, 92). Decisions at time t only ever see data <= t; future closes |
| are used strictly afterward, only to score the already-fixed decision |
| (section 53: "Only AFTER execution/decision is fixed may future data be |
| used to calculate realized PnL"). |
| |
| Fill-price convention (documented per section 25's "all assumptions |
| must be explicit, configurable, and auditable"): entry fills at the |
| CLOSE of the bar `execution_delay_bars` after the decision origin; exit |
| fills at the close `horizon` bars after that. Swap for next-bar-OPEN if |
| you want a stricter convention — it's one line, marked below. |
| |
| Single-position convention: this simulates ONE account with no margin |
| or leverage model, so a new trade is never opened while a previously |
| opened trade hasn't exited yet (decide() has no notion of position |
| state, so this is enforced here). Without this, a fresh full-notional |
| trade would be opened at every origin regardless of what's still open |
| -- during a run of same-direction signals that's several positions |
| stacked on the same equity curve at once, which silently multiplies |
| exposure the cost model and risk metrics know nothing about. |
| """ |
| from __future__ import annotations |
|
|
| from dataclasses import dataclass |
| from typing import List, Optional |
|
|
| import numpy as np |
| import pandas as pd |
|
|
| import config as cfg |
| from calibration import ProbabilityCalibrator |
| from decision_engine import Decision, decide |
| from features import warm_up_valid_from |
| from leakage_audit import LeakageError, audit_feature_frame |
|
|
|
|
| class InsufficientDataError(Exception): |
| """Raised when the feature frame has fewer bars than this backtest |
| configuration needs for even a single origin -> exit sequence. |
| Distinct from the ValueError warm_up_valid_from raises (which means |
| indicator warm-up is impossible at all): this means warm-up succeeded |
| but there isn't enough data left over for context_length plus at |
| least one tradeable origin. Previously this case wasn't checked for |
| at all -- feature_df.index[start_idx] would raise a bare, confusing |
| IndexError once start_idx ran past the end of the frame. Checking |
| explicitly means the message can say what failed, not just that |
| something did (spec section 98).""" |
|
|
|
|
| @dataclass |
| class Trade: |
| decision_ts: pd.Timestamp |
| entry_ts: pd.Timestamp |
| exit_ts: pd.Timestamp |
| action: str |
| entry_price: float |
| exit_price: float |
| forecast_return: float |
| gross_pnl: float |
| net_pnl: float |
| costs: float |
| exit_reason: str = "horizon" |
| prob_favorable: float = float("nan") |
| interval_width_pct: float = float("nan") |
|
|
|
|
| def _scan_for_early_exit(feature_df: pd.DataFrame, exec_idx: int, horizon_exit_idx: int, |
| action: str, entry_price: float, risk: Optional[cfg.RiskManagement]) -> tuple: |
| """First bar strictly after entry whose high/low breaches a stop-loss |
| or take-profit level, or horizon_exit_idx if neither ever triggers |
| (or risk management is off). Returns (exit_idx, exit_reason). |
| |
| The exit PRICE is still that bar's CLOSE, not the stop/TP level |
| itself -- consistent with the rest of this engine's documented |
| close-only fill convention (module docstring), and because intrabar |
| high/low only tells us a threshold was crossed, not what price you'd |
| actually have been filled at; inventing an exact fill price from |
| that would be manufacturing data we don't have (spec section 88). |
| This is conservative in both directions: a stop may report a |
| smaller loss than a real gap-through would produce, and a |
| take-profit may report less gain than the intrabar peak. If both |
| levels are breached within the same bar, we can't tell from OHLC |
| alone which was touched first, so this assumes the stop -- the |
| conservative read when direction of causation is unknown.""" |
| if risk is None or (risk.stop_loss_frac is None and risk.take_profit_frac is None): |
| return horizon_exit_idx, "horizon" |
|
|
| if action == "BUY": |
| stop_level = entry_price * (1 - risk.stop_loss_frac) if risk.stop_loss_frac else None |
| tp_level = entry_price * (1 + risk.take_profit_frac) if risk.take_profit_frac else None |
| else: |
| stop_level = entry_price * (1 + risk.stop_loss_frac) if risk.stop_loss_frac else None |
| tp_level = entry_price * (1 - risk.take_profit_frac) if risk.take_profit_frac else None |
|
|
| for j in range(exec_idx + 1, horizon_exit_idx + 1): |
| lo = float(feature_df["low"].iloc[j]) |
| hi = float(feature_df["high"].iloc[j]) |
| hit_stop = stop_level is not None and ( |
| (action == "BUY" and lo <= stop_level) or (action == "SELL" and hi >= stop_level) |
| ) |
| hit_tp = tp_level is not None and ( |
| (action == "BUY" and hi >= tp_level) or (action == "SELL" and lo <= tp_level) |
| ) |
| if hit_stop: |
| return j, "stop_loss" |
| if hit_tp: |
| return j, "take_profit" |
| return horizon_exit_idx, "horizon" |
|
|
|
|
| def _build_bar_by_bar_equity(feature_df: pd.DataFrame, start_idx: int, end_idx: int, |
| trades: List["Trade"]) -> pd.Series: |
| """TRUE mark-to-market equity: one point per BAR from start_idx to |
| end_idx inclusive (not one point per origin/trade, and never a gap |
| during an open position). While flat, equity is unchanged. While a |
| position is open, each bar's value reflects that bar's UNREALIZED |
| gross return relative to entry -- capturing genuine interim |
| drawdown/runup a trade-event-only curve hides entirely (a position |
| that spikes 2% against you mid-hold before recovering to a 0.3% |
| loss at exit previously showed as a single flat jump straight to |
| -0.3%, with the 2% excursion invisible to max_drawdown/Sharpe/ |
| Sortino/Calmar). The round-trip cost is applied once, at the |
| realized exit bar (matching net_pnl's own definition), not smeared |
| across intermediate marks -- you haven't paid the exit-side cost |
| until you've actually exited.""" |
| idx = feature_df.index[start_idx:end_idx + 1] |
| closes = feature_df["close"].iloc[start_idx:end_idx + 1].to_numpy() |
| n = len(idx) |
| equity = np.empty(n) |
|
|
| entry_map = {t.entry_ts: t for t in trades} |
| base = 1.0 |
| active = None |
| entry_base = 1.0 |
| for k in range(n): |
| ts = idx[k] |
| if active is None and ts in entry_map: |
| active = entry_map[ts] |
| entry_base = base |
| if active is not None: |
| close_k = float(closes[k]) |
| unreal = (close_k / active.entry_price - 1.0) if active.action == "BUY" \ |
| else (active.entry_price / close_k - 1.0) |
| if ts == active.exit_ts: |
| equity[k] = entry_base * (1.0 + active.net_pnl) |
| base = equity[k] |
| active = None |
| else: |
| equity[k] = entry_base * (1.0 + unreal) |
| else: |
| equity[k] = base |
| return pd.Series(equity, index=idx) |
|
|
|
|
| class WalkForwardBacktester: |
| """`forecaster` must expose .predict(context_df, target_col, |
| past_covariate_cols, prediction_length, freq) -> an object with |
| .quantiles / .quantile_levels (see moirai_model.QuantileForecast). |
| In production this is always a MoiraiMoEForecaster; tests may inject |
| a deterministic stub to verify the pipeline's causal behavior — |
| that stub is never imported by app.py (see tests/fixtures.py). |
| |
| This is the ONLY place that turns a Moirai forecast into a |
| simulated trade: app.py's Live tab calls decide() directly with the |
| same DecisionConfig/cost/calibrator a backtest run would use, so |
| there is exactly one decision pipeline, not two independently- |
| maintained ones (spec: "backtest profitable but live behaves |
| differently" is a structural risk this project treats as a real bug |
| class, not just a documentation gap).""" |
|
|
| def __init__(self, forecaster, timeframe: str, asset_class: str = "forex", |
| horizon: int = 5, context_length: int = 256, |
| decision_cfg: Optional[cfg.DecisionConfig] = None, |
| risk: Optional[cfg.RiskManagement] = None, |
| calibrator: Optional["ProbabilityCalibrator"] = None): |
| self.forecaster = forecaster |
| self.timeframe = timeframe |
| self.freq = cfg.TIMEFRAME_TO_PANDAS_FREQ[timeframe] |
| self.cost_profile = cfg.COST_PROFILES[asset_class] |
| self.horizon = horizon |
| self.context_length = context_length |
| self.decision_cfg = decision_cfg or cfg.DecisionConfig() |
| self.risk = risk |
| self.calibrator = calibrator |
|
|
| def run(self, feature_df: pd.DataFrame, start_idx: Optional[int] = None, step: int = 1) -> dict: |
| cost_frac = cfg.round_trip_cost_frac(self.cost_profile) |
| warm = warm_up_valid_from(feature_df) |
|
|
| n = len(feature_df) |
| delay = self.cost_profile.execution_delay_bars |
| min_start = max(start_idx or 0, warm, warm + self.context_length) |
| last_origin = n - self.horizon - delay - 1 |
| if min_start > last_origin: |
| bars_needed = min_start + self.horizon + delay + 1 |
| raise InsufficientDataError( |
| f"Not enough data for a single backtest origin: {n} bar(s) " |
| f"available, but this configuration needs at least " |
| f"{bars_needed} (warm-up {warm} + context_length " |
| f"{self.context_length} + horizon {self.horizon} + " |
| f"execution_delay {delay} + 1). Widen the History window " |
| f"by at least {bars_needed - n} bar(s), or reduce context " |
| f"length / horizon." |
| ) |
| start_idx = min_start |
|
|
| trades: List[Trade] = [] |
| decisions_log: List[dict] = [] |
| errors: List[dict] = [] |
| n_origins = 0 |
| open_until_idx = -1 |
| end_idx = start_idx |
|
|
| |
| |
| |
| diag = { |
| "signals_buy": 0, "signals_sell": 0, |
| "hold_wide_interval": 0, "hold_no_edge": 0, |
| "hold_low_confidence_buy": 0, "hold_low_confidence_sell": 0, |
| "model_errors": 0, "skipped_tail_insufficient_future_bars": 0, |
| "skipped_overlap_position_open": 0, |
| "trades_opened_buy": 0, "trades_opened_sell": 0, |
| "exit_horizon": 0, "exit_stop_loss": 0, "exit_take_profit": 0, |
| } |
|
|
| for i in range(start_idx, max(start_idx, last_origin), step): |
| n_origins += 1 |
| origin_ts = feature_df.index[i] |
| context = feature_df.iloc[i - self.context_length + 1: i + 1] |
|
|
| audit = audit_feature_frame(context, origin_ts) |
| if not audit["passed"]: |
| |
| |
| raise LeakageError(f"Backtest halted at {origin_ts}: {audit['violations']}") |
|
|
| try: |
| forecast = self.forecaster.predict( |
| context_df=context, target_col=cfg.TARGET_FEATURE, |
| past_covariate_cols=[c for c in cfg.PAST_DYNAMIC_FEATURES if c in context.columns], |
| prediction_length=self.horizon, freq=self.freq, |
| ) |
| except Exception as e: |
| errors.append({"origin_ts": str(origin_ts), "error": str(e), "error_type": type(e).__name__}) |
| diag["model_errors"] += 1 |
| continue |
|
|
| price_now = float(feature_df["close"].iloc[i]) |
| decision: Decision = decide(forecast, price_now, self.decision_cfg, |
| cost_estimate_frac=cost_frac, calibrator=self.calibrator) |
| decisions_log.append({"origin_ts": origin_ts, "action": decision.action, |
| "reject_stage": decision.reject_stage}) |
|
|
| if decision.action == "BUY": |
| diag["signals_buy"] += 1 |
| elif decision.action == "SELL": |
| diag["signals_sell"] += 1 |
| elif f"hold_{decision.reject_stage}" in diag: |
| diag[f"hold_{decision.reject_stage}"] += 1 |
| else: |
| |
| |
| |
| |
| |
| raise AssertionError( |
| f"decide() returned action={decision.action!r} with an " |
| f"unrecognized reject_stage={decision.reject_stage!r} — " |
| f"diagnostics cannot account for this origin" |
| ) |
|
|
| exec_idx = i + delay |
| exit_idx = i + delay + self.horizon |
| if exit_idx >= n: |
| diag["skipped_tail_insufficient_future_bars"] += 1 |
| continue |
| end_idx = max(end_idx, exit_idx) |
|
|
| if i < open_until_idx: |
| |
| |
| |
| |
| |
| |
| if decision.action in ("BUY", "SELL"): |
| diag["skipped_overlap_position_open"] += 1 |
| continue |
|
|
| entry_price = float(feature_df["close"].iloc[exec_idx]) |
|
|
| if decision.action == "HOLD": |
| continue |
|
|
| actual_exit_idx, exit_reason = _scan_for_early_exit( |
| feature_df, exec_idx, exit_idx, decision.action, entry_price, self.risk |
| ) |
| exit_price = float(feature_df["close"].iloc[actual_exit_idx]) |
|
|
| if decision.action == "BUY": |
| gross = exit_price / entry_price - 1.0 |
| else: |
| gross = entry_price / exit_price - 1.0 |
| net = gross - cost_frac |
|
|
| trades.append(Trade( |
| decision_ts=origin_ts, entry_ts=feature_df.index[exec_idx], |
| exit_ts=feature_df.index[actual_exit_idx], action=decision.action, |
| entry_price=entry_price, exit_price=exit_price, |
| forecast_return=decision.expected_return, gross_pnl=gross, net_pnl=net, costs=cost_frac, |
| exit_reason=exit_reason, prob_favorable=decision.prob_favorable, |
| interval_width_pct=decision.interval_width_pct, |
| )) |
| diag[f"trades_opened_{decision.action.lower()}"] += 1 |
| diag[f"exit_{exit_reason}"] += 1 |
| open_until_idx = actual_exit_idx |
|
|
| |
| |
| |
| |
| accounted = (diag["signals_buy"] + diag["signals_sell"] + diag["hold_wide_interval"] |
| + diag["hold_no_edge"] + diag["hold_low_confidence_buy"] + diag["hold_low_confidence_sell"] |
| + diag["model_errors"] + diag["skipped_tail_insufficient_future_bars"]) |
| |
| |
| |
| |
| assert accounted == n_origins, ( |
| f"diagnostic bucket mismatch: {accounted} accounted vs {n_origins} origins — " |
| f"diagnostics do not fully explain every origin's outcome" |
| ) |
|
|
| equity_curve = _build_bar_by_bar_equity(feature_df, start_idx, end_idx, trades) |
|
|
| return { |
| "trades": trades, |
| "decisions_log": decisions_log, |
| "equity_curve": equity_curve, |
| "num_origins": n_origins, |
| "errors": errors, |
| "diagnostics": diag, |
| "requested_risk": {"stop_loss_frac": self.risk.stop_loss_frac if self.risk else None, |
| "take_profit_frac": self.risk.take_profit_frac if self.risk else None}, |
| } |
|
|