SPB Git

spb/anomaly-atlas Public License

Systematic discovery & rigorous validation of statistical anomalies in open HF market data (hfmarketdata.io) — pre-registered, artifact-null-driven, fully reproducible. Live atlas: www.anomaly-atlas.io

Python 61.4% JavaScript 28.7% CSS 8.6% Shell 0.7% Makefile 0.5%
11.8 KB · 285 lines python
Raw Blame History
1# =============================================================================2#  Project   : anomaly-atlas3#  File      : experiments/micro/expG_cost_frontier/benchmark.py4#  Purpose   : Cost frontier: net-of-spread sweep over the double-filtered pool5#  Author    : Simon-Pierre Boucher6#  Contact   : contact@spboucher.ai7#  Data src  : hfmarketdata.io (sole data source)8#  Created   : 2026-08-129#  Modified  : 2026-08-1210#  Platform  : macOS / Apple Silicon (arm64)11#  License   : All rights reserved (research code)12# =============================================================================13"""Experiment G — the cost frontier (protocol pre-specified in hypothesis.md).1415Pool = mechanical intersection of committed expC/expD/expF outputs. For each16rule: gross daily stream + daily turnover, then net = gross − κ·(EDGE/2)·17turnover across the declared κ sweep, with the analytic breakeven κ*.18"""1920from __future__ import annotations2122import glob23import json24import sys25from collections import defaultdict26from datetime import UTC, datetime27from pathlib import Path2829import numpy as np3031REPO_ROOT = Path(__file__).resolve().parents[3]32sys.path.insert(0, str(REPO_ROOT / "benchmarks"))33sys.path.insert(0, str(REPO_ROOT / "src"))3435from hardware_manifest import collect_manifest  # noqa: E4023637from anomaly_atlas.data.cleaning import RTH_SLOTS, rth_day_grids  # noqa: E40238from anomaly_atlas.data.hf_client import HFMarketDataClient  # noqa: E40239from anomaly_atlas.data.universe import TRAIN_SUBPERIODS  # noqa: E40240from anomaly_atlas.stats.bootstrap import moving_block_bootstrap, percentile_ci  # noqa: E40241from anomaly_atlas.validation.artifacts import edge_spread  # noqa: E4024243ADJ = "adj_split"44KAPPAS = [0.0, 0.1, 0.25, 0.5, 1.0, 2.0]45ONE_MIN_WINDOW = ("2014-01-01", "2016-01-01")46D_WINDOW = ("2014-01-01", "2016-01-01")474849def latest(pattern: str) -> dict:50    return json.loads(Path(sorted(glob.glob(pattern))[-1]).read_text())515253def pool_from_committed_results() -> tuple[list[dict], list[dict]]:54    """Rebuild the double-filtered pool mechanically (no hand-picking)."""55    rc = latest(str(REPO_ROOT / "results/expC_reversion_scan/*/results.json"))56    rf = latest(str(REPO_ROOT / "results/expF_multiple_testing/*/results.json"))57    rd = latest(str(REPO_ROOT / "results/expD_leadlag_scan/*/results.json"))5859    cells = rc["cells"]60    for c in cells:61        c["vr30_excess"] = c["vr30"] - (1 + 2 * c["ac1"] * (1 - 1 / 30))62    triage = {(c["ticker"], c["timeframe"], c["period"]) for c in cells63              if c.get("fdr_vr30") and c["vr30_excess"] < -0.05}64    spa_r = {tuple(s[3:].split(":")) for s in rf["funnel"]["spa_step1_survivors"]65             if s[1] == "R"}  # (ticker, timeframe)66    rev_pool = [{"ticker": t, "timeframe": tf, "period": per,67                 "asset": "etf" if t in ("SPY", "QQQ") else "stock"}68                for (t, tf, per) in sorted(triage) if (t, tf) in spa_r]6970    ll_pool = []71    for p in rd["pairs"]:72        if p["window"] != "2014-2015" or p["bucket"] == "index":73            continue74        if not (p.get("fdr_fresh_+1") or p.get("fdr_fresh_-1")):75            continue76        follower = p["pair"].split("->")[1] if p["pair"].startswith("SPY->") else "SPY"77        leader = "SPY" if follower != "SPY" else p["pair"].split("->")[0].split("[")[0]78        sign = np.sign(p["fresh_xcorr"]["1"]) if p.get("fdr_fresh_+1") else np.sign(79            p["fresh_xcorr"]["-1"])80        ll_pool.append({"pair": p["pair"], "leader": leader, "follower": follower,81                        "sign": int(sign), "bucket": p["bucket"]})82    return rev_pool, ll_pool838485def contrarian_stream(bars: list[dict], timeframe: str) -> dict[str, tuple[float, float]]:86    """day -> (gross rule return, turnover) for the +contrarian rule."""87    by_day: dict[str, list[float]] = defaultdict(list)88    for b in bars:89        dt = b["datetime"]90        if timeframe == "1day" or "09:30" <= dt[11:16] < "16:00":91            by_day[dt[:10]].append(np.log(b["close"]))92    days = sorted(by_day)93    out: dict[str, tuple[float, float]] = {}94    if timeframe == "1day":95        closes = np.array([by_day[d][0] for d in days])96        r = np.diff(closes)97        pos_prev = 0.098        for i in range(1, len(r)):99            pos = -np.sign(r[i - 1])100            out[days[i + 1]] = (float(pos * r[i]), float(abs(pos - pos_prev)))101            pos_prev = pos102        return out103    for d in days:104        p = np.array(by_day[d])105        if len(p) < 3:106            continue107        r = np.diff(p)108        pos = -np.sign(r[:-1])109        gross = float(np.sum(pos * r[1:]))110        turnover = float(abs(pos[0]) + np.sum(np.abs(np.diff(pos))) + abs(pos[-1]))111        out[d] = (gross, turnover)112    return out113114115def leadlag_stream(gx: dict, gy: dict, day_list: list[str], sign: int116                   ) -> dict[str, tuple[float, float]]:117    """day -> (gross, turnover) for pos_t = sign * sign(x_{t-1}) on fresh minutes."""118    out: dict[str, tuple[float, float]] = {}119    for day in day_list:120        px, py = gx.get(day), gy.get(day)121        if px is None or py is None:122            continue123        ox, oy = np.isfinite(px), np.isfinite(py)124        fpx, fpy = px.copy(), py.copy()125        for t in range(1, RTH_SLOTS):126            if not ox[t]:127                fpx[t] = fpx[t - 1]128            if not oy[t]:129                fpy[t] = fpy[t - 1]130        rx, ry = np.diff(fpx), np.diff(fpy)131        fx = ox[1:] & ox[:-1]132        fy = oy[1:] & oy[:-1]133        keep = fx[:-1] & fy[1:] & np.isfinite(rx[:-1]) & np.isfinite(ry[1:])134        if keep.sum() < 30:135            continue136        pos = np.zeros(len(ry))137        pos[1:][keep] = sign * np.sign(rx[:-1][keep])138        gross = float(np.sum(pos[1:] * ry[1:]))139        turnover = float(np.sum(np.abs(np.diff(np.concatenate([[0.0], pos, [0.0]])))))140        out[day] = (gross, turnover)141    return out142143144def half_spread(client: HFMarketDataClient, asset: str, ticker: str,145                adj: str | None, start: str, end: str) -> float:146    bars = client.get_bars(asset, ticker, "1day", adj, start, end)147    if len(bars) < 100:148        return float("nan")149    o = np.array([b["open"] for b in bars])150    h = np.array([b["high"] for b in bars])151    lo = np.array([b["low"] for b in bars])152    c = np.array([b["close"] for b in bars])153    s = edge_spread(o, h, lo, c)154    return s / 2.0 if np.isfinite(s) else float("nan")155156157def sanitize(obj):158    """Replace NaN/inf with None recursively — Python json emits bare NaN,159    which is invalid JSON for every other consumer (incl. the web platform)."""160    if isinstance(obj, dict):161        return {k: sanitize(v) for k, v in obj.items()}162    if isinstance(obj, list):163        return [sanitize(v) for v in obj]164    if isinstance(obj, float) and not np.isfinite(obj):165        return None166    return obj167168169def sweep(stream: dict[str, tuple[float, float]], hs: float) -> dict:170    days = sorted(d for d in stream if np.isfinite(stream[d][0]) and np.isfinite(stream[d][1]))171    gross = np.array([stream[d][0] for d in days])172    turn = np.array([stream[d][1] for d in days])173    res = {174        "n_days": len(days),175        "gross_mean_daily_bp": round(float(gross.mean()) * 1e4, 3),176        "turnover_per_day": round(float(turn.mean()), 2),177        "half_spread_bp": round(hs * 1e4, 3) if np.isfinite(hs) else None,178        "net": {},179    }180    if not np.isfinite(hs) or hs <= 0:181        res["kappa_star"] = None182        return res183    for k in KAPPAS:184        net = gross - k * hs * turn185        boot = moving_block_bootstrap(net, lambda x: float(np.mean(x)),186                                      block=21, n_boot=300, seed=42)187        lo_ci, hi_ci = percentile_ci(boot)188        res["net"][str(k)] = {189            "mean_daily_bp": round(float(net.mean()) * 1e4, 3),190            "ci95_bp": [round(lo_ci * 1e4, 3), round(hi_ci * 1e4, 3)],191            "positive": bool(net.mean() > 0),192        }193    denom = hs * turn.mean()194    res["kappa_star"] = round(float(gross.mean() / denom), 4) if denom > 0 else None195    return res196197198def main() -> None:199    run_utc = datetime.now(UTC)200    client = HFMarketDataClient()201    rev_pool, ll_pool = pool_from_committed_results()202    print("reversion pool:", [(r["ticker"], r["timeframe"], r["period"]) for r in rev_pool])203    print("leadlag pool:", [(p["pair"], p["sign"]) for p in ll_pool])204205    items = []206    for r in rev_pool:207        if r["period"] in TRAIN_SUBPERIODS:208            s, e = TRAIN_SUBPERIODS[r["period"]]209        else:210            s, e = ONE_MIN_WINDOW211        bars = client.get_bars(r["asset"], r["ticker"], r["timeframe"], ADJ, s, e)212        stream = contrarian_stream(bars, r["timeframe"])213        hs = half_spread(client, r["asset"], r["ticker"], ADJ, s, e)214        item = {"rule": f"R:{r['ticker']}:{r['timeframe']}:{r['period']}",215                "family": "reversion"} | sweep(stream, hs)216        items.append(item)217        print(item["rule"], "kappa* =", item["kappa_star"])218219    # lead-lag pool (2014-2015 window)220    s, e = D_WINDOW221    needed = {"SPY"} | {p["leader"] for p in ll_pool} | {p["follower"] for p in ll_pool}222    grids = {}223    for name in sorted(needed):224        if name == "ES":225            g = rth_day_grids(client.get_bars("futures", "ES", "1min",226                                              "contin_adj_ratio", s, e))227        else:228            asset = "etf" if name in ("SPY", "QQQ", "XLF", "XLE", "XLK", "XLV", "XLI",229                                      "XLY", "XLP", "XLU", "XLB") else "stock"230            g = rth_day_grids(client.get_bars(asset, name, "1min", ADJ, s, e))231        grids[name] = g232    day_list = sorted(grids["SPY"].keys())233    for p in ll_pool:234        stream = leadlag_stream(grids[p["leader"]], grids[p["follower"]],235                                day_list, p["sign"])236        traded = p["follower"]237        asset = ("futures" if traded == "ES" else238                 "etf" if traded in ("SPY", "QQQ") or traded.startswith("XL") else "stock")239        adj = "contin_adj_ratio" if traded == "ES" else ADJ240        hs = half_spread(client, asset, traded, adj, s, e)241        item = {"rule": f"L:{p['pair']}:{'+' if p['sign'] > 0 else '-'}",242                "family": "leadlag", "traded": traded} | sweep(stream, hs)243        items.append(item)244        print(item["rule"], "kappa* =", item["kappa_star"])245246    ks = [i["kappa_star"] for i in items247          if i["kappa_star"] is not None and np.isfinite(i["kappa_star"])]248    summary = {249        "pool_size": len(items),250        "n_kappa_star_defined": len(ks),251        "kappa_star_median": round(float(np.median(ks)), 4) if ks else None,252        "kappa_star_max": round(float(max(ks)), 4) if ks else None,253        "survivors_at": {str(k): [i["rule"] for i in items254                                  if i["net"].get(str(k), {}).get("positive")]255                         for k in (0.1, 0.25, 0.5, 1.0)},256        "intraday_survivor_at_1.0": [257            i["rule"] for i in items258            if ":1day:" not in i["rule"] and i["net"].get("1.0", {}).get("positive")],259    }260261    results = {262        "experiment": "expG_cost_frontier",263        "run_utc": run_utc.isoformat(),264        "author": "Simon-Pierre Boucher",265        "contact": "contact@spboucher.ai",266        "data_source": "hfmarketdata.io",267        "confidence_level": 0,268        "protocol": {"kappas": KAPPAS, "cost_model": "net = gross - k*(EDGE/2)*turnover",269                     "pool": "mechanical intersection of committed expC/expD/expF results"},270        "items": items,271        "summary": summary,272        "client_stats": vars(client.stats) | {"refreshes": list(client.stats.refreshes)},273        "manifest": collect_manifest(),274    }275    out_dir = REPO_ROOT / "results" / "expG_cost_frontier" / run_utc.strftime("%Y%m%dT%H%M%SZ")276    out_dir.mkdir(parents=True)277    (out_dir / "results.json").write_text(278        json.dumps(sanitize(results), indent=2, allow_nan=False) + "\n")279    print(f"\nwrote {out_dir.relative_to(REPO_ROOT)}/results.json")280    print(json.dumps(summary, indent=1))281282283if __name__ == "__main__":284    main()285