Lk / backtest_engine.py
Kashaf1's picture
Upload moirai_forecast_app contents
ef20ebe
Raw
History Blame Contribute Delete
17.2 kB
"""
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" # "horizon" | "stop_loss" | "take_profit"
prob_favorable: float = float("nan") # decision.prob_favorable at open, for post-hoc entry-filter analysis
interval_width_pct: float = float("nan") # decision.interval_width_pct at open, same reason
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 # None => no early exit, identical to every pre-existing backtest
self.calibrator = calibrator # None => raw probabilities, identical to every pre-existing backtest
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) # raises ValueError if warm-up is impossible at all
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 # bar index the currently-open trade exits at; -1 == flat
end_idx = start_idx # highest bar index any origin's exit reaches; grown below
# Diagnostics (spec section 98 / audit item 18): every origin ends
# up in exactly one of these buckets, so they must sum to
# n_origins -- asserted at the end rather than just hoped for.
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"]:
# Fail loudly rather than silently dropping the bad origin
# (spec section 26).
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:
# Should be unreachable -- decide() only ever sets
# reject_stage to one of the 4 keys above when action is
# HOLD. Surfacing this loudly rather than silently
# dropping the origin from the count is the whole point
# of asserting `accounted == n_origins` below.
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:
# Still holding a position opened at an earlier origin --
# don't stack a second full-notional trade on top of it
# (see "Single-position convention" in the module
# docstring). The already-open trade's own exit will
# still post its P&L to the equity curve when its turn
# comes; this origin just isn't tradable right now.
if decision.action in ("BUY", "SELL"):
diag["skipped_overlap_position_open"] += 1
continue
entry_price = float(feature_df["close"].iloc[exec_idx]) # fill convention, see module docstring
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]) # future data — used only to SCORE, decision already fixed
if decision.action == "BUY":
gross = exit_price / entry_price - 1.0
else: # SELL
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
# Every origin must land in exactly one bucket -- a mismatch here
# would mean the diagnostics themselves are lying, which is worse
# than not having them (audit item 18 exists so counts can be
# trusted, not just displayed).
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"])
# (BUY/SELL signals still count once here even when later skipped
# for overlap -- skipped_overlap_position_open is a sub-count of
# signals_buy/signals_sell, not an additional bucket, since the
# decision itself was still a real BUY/SELL signal.)
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},
}