Brings the durable artifacts of research/portfolio-capacity-rebalancing onto main so the rationale for raising the count cap lives with the code that cites it. The matrix runner, the research simulator hooks and the study's unit tests are deliberately left behind; they remain at tag research/portfolio-capacity-final. Corrects conclusions that were reached on EV per trade and are now superseded: the findings doc's decisions 1 (keep cap 10) and 4 (run the risk-floor A/B) are struck through and answered in a new correction section, and the research README and phase-A matrix entries are updated to match. The frozen specification itself is untouched -- its recorded SHA-256 f1e37783 still verifies. effective-risk-floor-ab.md is retained but marked CLOSED/NEGATIVE: the study it proposes is already answered by cap15 vs cash_unbounded (-0.753pp CAGR while EV/trade rises), and its EV-based pass rule would have shipped it. scripts/research_rankings.py replaces a fourth copy of the historical rank-map helper; run_research_matrix, run_execution_recovery_matrix and run_daily_reentry_matrix now share it. The shared version adds a duplicate observation guard and a deterministic symbol tie-break the copies lacked. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
440 lines
17 KiB
Python
440 lines
17 KiB
Python
"""Execution-recovery matrix: is the close→next_open gap recoverable by scheduling?
|
||
|
||
Hypothesis (from Phase A A4)
|
||
----------------------------
|
||
The 1.77 → 1.20 full-period Sharpe gap under next_open fill is mostly overnight
|
||
momentum drift that a 07:00-Berlin scanner cannot earn. Near-close / MOC-style
|
||
execution (scan ~15:45 ET, fill at/near that close) should recover it.
|
||
|
||
Pre-registered arms (N for DSR = 4)
|
||
----------------------------------
|
||
1. ``close_control`` — historical optimistic control (signal = fill at same close).
|
||
2. ``next_open`` — honest overnight scanner (decision baseline for future promotion).
|
||
3. ``stale_close`` — signal at t−1 close, fill at t close (one-session-stale MOC proxy).
|
||
Expectation: ≈ close_control; if so, the gap is scheduling, not physics.
|
||
4. ``next_open_gap2`` — next_open but skip entries that open > +2% above signal close.
|
||
Measure whether large gap-ups are toxic or the best continuations.
|
||
|
||
Promotion / read rule (pre-registered)
|
||
--------------------------------------
|
||
- Decision baseline for *future* strategy work: ``next_open``.
|
||
- Recovery success for ``stale_close``: validation Sharpe within 0.5×SE of
|
||
``close_control`` **and** validation Sharpe ≥ ``next_open``; train Sharpe not
|
||
worse than close_control by more than 0.5×SE. State 1-SE distinguishability.
|
||
- ``next_open_gap2`` is measurement-only vs ``next_open`` (no auto-promote to live).
|
||
|
||
Reuses the same daily candidate cache as the Phase A matrix when the cache key
|
||
matches.
|
||
|
||
Usage
|
||
-----
|
||
python scripts/run_execution_recovery_matrix.py backtest_snapshots/prod.sqlite \\
|
||
--workers 7 --allow-spawn \\
|
||
--candidate-cache reports/.cache/research-cands.pkl \\
|
||
--out reports/execution-recovery-matrix.json
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import asyncio
|
||
import json
|
||
import multiprocessing
|
||
import os
|
||
import pickle
|
||
import sys
|
||
from concurrent.futures import ProcessPoolExecutor, as_completed
|
||
from datetime import date, datetime
|
||
from pathlib import Path
|
||
from typing import Any
|
||
|
||
from sqlalchemy import select
|
||
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine
|
||
|
||
ROOT = Path(__file__).resolve().parents[1]
|
||
if str(ROOT) not in sys.path:
|
||
sys.path.insert(0, str(ROOT))
|
||
|
||
from scripts.research_rankings import ( # noqa: E402
|
||
_live_universe_rank_map,
|
||
_period_percentiles,
|
||
)
|
||
|
||
# Must match Phase A cache when reusing research-cands.pkl
|
||
CACHE_VERSION = "research-matrix-v1-daily-prod"
|
||
|
||
PRE_REGISTERED_ARMS: tuple[dict[str, Any], ...] = (
|
||
{
|
||
"id": "close_control",
|
||
"label": "Close fill (historical control)",
|
||
"fill_mode": "close",
|
||
},
|
||
{
|
||
"id": "next_open",
|
||
"label": "Next-open fill (decision baseline)",
|
||
"fill_mode": "next_open",
|
||
},
|
||
{
|
||
"id": "stale_close",
|
||
"label": "Stale-signal close fill (MOC proxy: signal t-1, fill t close)",
|
||
"fill_mode": "stale_close",
|
||
},
|
||
{
|
||
"id": "next_open_gap2",
|
||
"label": "Next-open + skip gap-up > 2%",
|
||
"fill_mode": "next_open",
|
||
"max_entry_gap_pct": 0.02,
|
||
},
|
||
)
|
||
N_TRIALS = len(PRE_REGISTERED_ARMS)
|
||
|
||
|
||
def _sqlite_url(path: Path) -> str:
|
||
return f"sqlite+aiosqlite:///{path.resolve().as_posix()}"
|
||
|
||
|
||
def _parse_args() -> argparse.Namespace:
|
||
p = argparse.ArgumentParser(
|
||
description=__doc__,
|
||
formatter_class=argparse.RawDescriptionHelpFormatter,
|
||
)
|
||
p.add_argument("snapshot")
|
||
p.add_argument("--workers", type=int, default=6)
|
||
p.add_argument("--allow-spawn", action="store_true")
|
||
p.add_argument("--out", default=None)
|
||
p.add_argument("--candidate-cache", default=None)
|
||
p.add_argument("--validation-split", default="2024-07-01")
|
||
p.add_argument("--cadence", choices=("daily", "weekly"), default="daily")
|
||
p.add_argument("--quiet", action="store_true")
|
||
return p.parse_args()
|
||
|
||
|
||
def _window(arm: dict, name: str) -> dict | None:
|
||
for row in arm.get("windows") or []:
|
||
if row.get("window") == name:
|
||
return row
|
||
return None
|
||
|
||
|
||
def _grade_stale_close(close_arm: dict, next_arm: dict, stale_arm: dict) -> dict:
|
||
c_val = _window(close_arm, "validation") or {}
|
||
n_val = _window(next_arm, "validation") or {}
|
||
s_val = _window(stale_arm, "validation") or {}
|
||
c_tr = _window(close_arm, "train") or {}
|
||
s_tr = _window(stale_arm, "train") or {}
|
||
keys = ("sharpe", "sharpe_se")
|
||
if any(c_val.get(k) is None for k in keys) or s_val.get("sharpe") is None:
|
||
return {"recover": False, "reason": "missing Sharpe rows"}
|
||
se = float(c_val.get("sharpe_se") or s_val.get("sharpe_se") or 0.0)
|
||
half_se = 0.5 * se if se > 0 else 0.0
|
||
cs, ss, ns = float(c_val["sharpe"]), float(s_val["sharpe"]), n_val.get("sharpe")
|
||
cts, sts = c_tr.get("sharpe"), s_tr.get("sharpe")
|
||
near_close = abs(ss - cs) <= half_se if half_se > 0 else abs(ss - cs) < 0.05
|
||
beats_next = ns is None or ss >= float(ns)
|
||
train_ok = (
|
||
cts is None
|
||
or sts is None
|
||
or float(sts) >= float(cts) - half_se
|
||
)
|
||
recover = near_close and beats_next and train_ok
|
||
return {
|
||
"recover": recover,
|
||
"near_close_control": near_close,
|
||
"beats_next_open": beats_next,
|
||
"train_ok": train_ok,
|
||
"validation_delta_vs_close": round(ss - cs, 4),
|
||
"validation_delta_vs_next_open": (
|
||
round(ss - float(ns), 4) if ns is not None else None
|
||
),
|
||
"half_se": half_se,
|
||
"reason": (
|
||
"stale_close recovers close-fill economics (within 0.5 SE) and beats next_open"
|
||
if recover
|
||
else "stale_close does not meet recovery criteria — see flags"
|
||
),
|
||
}
|
||
|
||
|
||
def _markdown(report: dict) -> str:
|
||
lines = [
|
||
f"# Execution recovery matrix — {report.get('generated_at', '')}",
|
||
"",
|
||
f"Validation split: **{report.get('validation_split')}**. N for DSR: **{report.get('n_trials')}**.",
|
||
"",
|
||
"| arm | window | Sharpe | SE | CAGR | MaxDD | trades | gap/drift |",
|
||
"|---|---|---:|---:|---:|---:|---:|---|",
|
||
]
|
||
for arm in report.get("arms") or []:
|
||
for w in arm.get("windows") or []:
|
||
slip = w.get("overnight_slippage") or w.get("signal_to_fill_drift") or {}
|
||
slip_s = (
|
||
f"mean {slip.get('mean_pct')}% n={slip.get('n')}"
|
||
if slip
|
||
else "—"
|
||
)
|
||
if w.get("skipped_gap_cap") is not None:
|
||
slip_s += f"; gap_skips={w.get('skipped_gap_cap')}"
|
||
lines.append(
|
||
f"| {arm.get('id')} | {w.get('window')} | {w.get('sharpe')} | "
|
||
f"{w.get('sharpe_se')} | {w.get('cagr_pct')} | {w.get('max_drawdown_pct')} | "
|
||
f"{w.get('trades')} | {slip_s} |"
|
||
)
|
||
rec = report.get("recovery") or {}
|
||
lines.extend(
|
||
[
|
||
"",
|
||
"## Recovery decision (stale_close)",
|
||
"",
|
||
f"- **recover: {rec.get('recover')}** — {rec.get('reason')}",
|
||
f"- flags: { {k: rec.get(k) for k in ('near_close_control', 'beats_next_open', 'train_ok', 'validation_delta_vs_close', 'validation_delta_vs_next_open', 'half_se')} }",
|
||
"",
|
||
]
|
||
)
|
||
return "\n".join(lines)
|
||
|
||
|
||
def _checkpoint(path: Path, report: dict) -> None:
|
||
path.parent.mkdir(parents=True, exist_ok=True)
|
||
tmp = path.with_suffix(path.suffix + ".tmp")
|
||
tmp.write_text(json.dumps(report, indent=2, default=str), encoding="utf-8")
|
||
tmp.replace(path)
|
||
path.with_suffix(".md").write_text(_markdown(report), encoding="utf-8")
|
||
|
||
|
||
async def _main() -> None:
|
||
args = _parse_args()
|
||
snapshot = Path(args.snapshot)
|
||
if not snapshot.exists():
|
||
raise SystemExit(f"Snapshot not found: {snapshot}")
|
||
validation_split = date.fromisoformat(args.validation_split)
|
||
out_path = (
|
||
Path(args.out)
|
||
if args.out
|
||
else Path("reports")
|
||
/ f"execution-recovery-matrix-{datetime.now().strftime('%Y%m%d-%H%M%S')}.json"
|
||
)
|
||
|
||
os.environ["BACKTEST_SNAPSHOT_OFFLINE"] = "1"
|
||
if args.allow_spawn:
|
||
os.environ["BACKTEST_ALLOW_SPAWN"] = "1"
|
||
|
||
from app.models.ticker import Ticker
|
||
from app.services import backtest_service as bt
|
||
from app.services.admin_service import get_activation_config
|
||
from app.services.paper_trade_service import get_exit_policy
|
||
from app.services.recommendation_service import get_recommendation_config
|
||
|
||
db_engine = create_async_engine(_sqlite_url(snapshot), pool_pre_ping=True)
|
||
Session = async_sessionmaker(db_engine, class_=AsyncSession, expire_on_commit=False)
|
||
try:
|
||
async with Session() as db:
|
||
recommendation_config = await get_recommendation_config(db)
|
||
activation = await get_activation_config(db)
|
||
exit_config = await get_exit_policy(db)
|
||
benchmark_closes = await bt._load_benchmark_closes_for_backtest(
|
||
db, days=None, refresh=False
|
||
)
|
||
ticker_result = await db.execute(select(Ticker).order_by(Ticker.symbol))
|
||
symbols = [t.symbol for t in ticker_result.scalars().all()]
|
||
prices: dict[str, tuple] = {}
|
||
for index, symbol in enumerate(symbols, 1):
|
||
columns = await bt._fetch_columns(db, symbol)
|
||
if columns is not None:
|
||
prices[symbol] = columns
|
||
if not args.quiet and index % 50 == 0:
|
||
print(f"loaded prices: {index}/{len(symbols)}", flush=True)
|
||
finally:
|
||
await db_engine.dispose()
|
||
|
||
snapshot_stat = snapshot.stat()
|
||
cache_key = {
|
||
"version": CACHE_VERSION,
|
||
"snapshot": str(snapshot.resolve()),
|
||
"snapshot_size": snapshot_stat.st_size,
|
||
"snapshot_mtime_ns": snapshot_stat.st_mtime_ns,
|
||
"cadence": args.cadence,
|
||
"target_model": "production_gtl",
|
||
}
|
||
cache_path = Path(args.candidate_cache) if args.candidate_cache else None
|
||
qualified: list[dict] | None = None
|
||
entry_candidate_count = 0
|
||
|
||
if cache_path is not None and cache_path.exists():
|
||
with cache_path.open("rb") as handle:
|
||
cached = pickle.load(handle) # noqa: S301
|
||
if cached.get("key") == cache_key:
|
||
qualified = list(cached["qualified_candidates"])
|
||
entry_candidate_count = int(cached.get("entry_candidate_count") or 0)
|
||
if not args.quiet:
|
||
print(f"loaded candidate cache: {cache_path}", flush=True)
|
||
|
||
if qualified is None:
|
||
workers = max(1, min(int(args.workers), max(1, multiprocessing.cpu_count() - 1)))
|
||
context = bt._mp_context() or multiprocessing.get_context("spawn")
|
||
replay_rows: list[dict] = []
|
||
with ProcessPoolExecutor(max_workers=workers, mp_context=context) as pool:
|
||
futures = {
|
||
pool.submit(
|
||
bt._replay_candidates_for_period,
|
||
symbol,
|
||
columns,
|
||
recommendation_config,
|
||
activation,
|
||
benchmark_closes,
|
||
date(1900, 1, 1),
|
||
args.cadence,
|
||
True,
|
||
True,
|
||
): symbol
|
||
for symbol, columns in prices.items()
|
||
}
|
||
for index, future in enumerate(as_completed(futures), 1):
|
||
replay_rows.extend(future.result())
|
||
if not args.quiet and index % 25 == 0:
|
||
print(f"replay: {index}/{len(futures)}", flush=True)
|
||
setup_candidates = [row for row in replay_rows if not row.get("_rank_only")]
|
||
rank_observations = [
|
||
row for row in replay_rows if row.get("_universe_rank_observation")
|
||
]
|
||
entry_candidate_count = len(setup_candidates)
|
||
live_ranks = _live_universe_rank_map(
|
||
rank_observations,
|
||
benchmark_closes,
|
||
bt.STRATEGY_RANK_MOMENTUM_WEIGHT,
|
||
)
|
||
threshold = float(activation.get("min_momentum_percentile", 80.0))
|
||
qualified = []
|
||
for setup in setup_candidates:
|
||
if setup.get("direction") != "long":
|
||
continue
|
||
candidate = {
|
||
k: v
|
||
for k, v in setup.items()
|
||
if not k.startswith("_universe_")
|
||
}
|
||
rank = live_ranks.get((str(setup["symbol"]), str(setup["date"])))
|
||
if rank is None:
|
||
continue
|
||
candidate[bt.PRODUCTION_PERCENTILE_KEY] = rank["momentum_percentile"]
|
||
candidate[bt.VOL_PERCENTILE_KEY] = rank["volatility_percentile"]
|
||
candidate[bt.RESIDUAL_HIGH_VOL_BLEND_80_20_KEY] = rank["strategy_rank"]
|
||
candidate["qualified"] = bt._momentum_qualifies(candidate, threshold)
|
||
if candidate["qualified"]:
|
||
qualified.append(candidate)
|
||
if cache_path is not None:
|
||
cache_path.parent.mkdir(parents=True, exist_ok=True)
|
||
with cache_path.open("wb") as handle:
|
||
pickle.dump(
|
||
{
|
||
"key": cache_key,
|
||
"entry_candidate_count": entry_candidate_count,
|
||
"qualified_candidates": qualified,
|
||
},
|
||
handle,
|
||
protocol=pickle.HIGHEST_PROTOCOL,
|
||
)
|
||
|
||
if not qualified:
|
||
raise SystemExit("No qualified candidates")
|
||
|
||
strategy = next(s for s in bt.PORTFOLIO_MONITOR_STRATEGIES if s.get("is_production"))
|
||
entry_config = bt._entry_variant_config(str(strategy["entry_variant"]))
|
||
assert entry_config is not None
|
||
ranking_key = str(entry_config.get("ranking_key") or entry_config["percentile_key"])
|
||
exit_policy = bt.LIVE_EXIT_MODE_TO_SIM.get(
|
||
str(exit_config.get("mode", "atr_trailing")), "atr_trail3"
|
||
)
|
||
hold_days = int(exit_config.get("hold_days", 30))
|
||
trail_multiplier = float(exit_config.get("atr_multiplier", bt.ATR_TRAIL_MULTIPLIER))
|
||
risk_per_trade = float(entry_config["risk_per_trade"])
|
||
max_positions = int(entry_config["max_positions"])
|
||
post_stop = bt._make_gate_reset_reentry_fn(
|
||
qualified, prices, cadence=args.cadence, ranking_key=ranking_key
|
||
)
|
||
|
||
report: dict[str, Any] = {
|
||
"generated_at": datetime.now().isoformat(),
|
||
"snapshot": str(snapshot.resolve()),
|
||
"validation_split": validation_split.isoformat(),
|
||
"n_trials": N_TRIALS,
|
||
"hypothesis": (
|
||
"stale_close (signal t-1, fill t close) recovers close-fill economics; "
|
||
"the next_open haircut is scheduling, not lost edge"
|
||
),
|
||
"decision_baseline": "next_open",
|
||
"qualified_longs": len(qualified),
|
||
"arms": [],
|
||
"recovery": {},
|
||
}
|
||
_checkpoint(out_path, report)
|
||
|
||
by_id: dict[str, dict] = {}
|
||
for arm in PRE_REGISTERED_ARMS:
|
||
if not args.quiet:
|
||
print(f"running {arm['id']} ...", flush=True)
|
||
windows = []
|
||
for window_name, start, end in (
|
||
("train", None, validation_split),
|
||
("validation", validation_split, None),
|
||
("full", None, None),
|
||
):
|
||
sim = bt._simulate_portfolio(
|
||
qualified,
|
||
prices,
|
||
benchmark_closes,
|
||
exit_policy,
|
||
hold_days,
|
||
ranking_key=ranking_key,
|
||
max_positions=max_positions,
|
||
risk_per_trade=risk_per_trade,
|
||
atr_trail_multiplier=trail_multiplier,
|
||
post_stop_reentry_fn=post_stop,
|
||
start_date=start,
|
||
end_date=end,
|
||
fill_mode=str(arm["fill_mode"]),
|
||
max_entry_gap_pct=arm.get("max_entry_gap_pct"),
|
||
include_trades=True,
|
||
)
|
||
if sim is None:
|
||
windows.append({"window": window_name, "error": "no_trades"})
|
||
continue
|
||
dsr = bt.deflated_sharpe_ratio(
|
||
sim.get("sharpe"),
|
||
sim.get("sharpe_se"),
|
||
N_TRIALS,
|
||
n_returns=sim.get("n_returns"),
|
||
return_skew=sim.get("return_skew"),
|
||
return_kurtosis=sim.get("return_kurtosis"),
|
||
)
|
||
sim.pop("trade_details", None)
|
||
sim.pop("equity_curve", None)
|
||
sim.pop("benchmark_curve", None)
|
||
sim.pop("reentry_events", None)
|
||
windows.append({"window": window_name, "dsr": dsr, **sim})
|
||
row = {"id": arm["id"], "label": arm["label"], "config": {
|
||
k: arm[k] for k in arm if k not in {"id", "label"}
|
||
}, "windows": windows}
|
||
by_id[arm["id"]] = row
|
||
report["arms"].append(row)
|
||
_checkpoint(out_path, report)
|
||
if not args.quiet:
|
||
val = _window(row, "validation") or {}
|
||
print(
|
||
f" {arm['id']}: val Sharpe={val.get('sharpe')} "
|
||
f"DD={val.get('max_drawdown_pct')} trades={val.get('trades')}",
|
||
flush=True,
|
||
)
|
||
|
||
if all(k in by_id for k in ("close_control", "next_open", "stale_close")):
|
||
report["recovery"] = _grade_stale_close(
|
||
by_id["close_control"], by_id["next_open"], by_id["stale_close"]
|
||
)
|
||
_checkpoint(out_path, report)
|
||
if not args.quiet:
|
||
print(f"wrote {out_path}", flush=True)
|
||
print(f"recovery: {report.get('recovery')}", flush=True)
|
||
|
||
|
||
if __name__ == "__main__":
|
||
asyncio.run(_main())
|