research: clean up closed Tier-1 scaffolding from branch

Drop intermediate history-depth reports, sector-residual runners/map/code hooks
(evidence stays in final reports + docs), and slim MacBook helper to ssl/earnings/
prod-book-matrix only. SSL bootstrap and archived research conclusions retained.
This commit is contained in:
2026-07-19 14:41:52 +02:00
parent 1c38a94dd0
commit bb8aa655a1
22 changed files with 88 additions and 5450 deletions
-233
View File
@@ -1,233 +0,0 @@
"""Build a local ticker → GICS sector map for research residualization.
Sources (in order):
1. Public S&P 500 constituents CSV (datasets/s-and-p-500-companies) — bulk, free.
2. Existing map file (resume).
3. FMP stable ``profile`` for still-missing symbols (budget ~250 req/day).
Writes ``data/research/ticker_sector_map.json``. Never touches production Postgres.
Example
-------
python scripts/build_ticker_sector_map.py \\
--snapshot backtest_snapshots/prod.sqlite
python scripts/build_ticker_sector_map.py --fmp-limit 50
"""
from __future__ import annotations
import argparse
import asyncio
import csv
import io
import json
import sys
import time
from datetime import datetime, timezone
from pathlib import Path
import httpx
from sqlalchemy import create_engine, text
ROOT = Path(__file__).resolve().parents[1]
if str(ROOT) not in sys.path:
sys.path.insert(0, str(ROOT))
from app.ssl_bootstrap import bootstrap_ssl # noqa: E402
bootstrap_ssl()
from app.services.sector_map import ( # noqa: E402
DEFAULT_SECTOR_MAP_PATH,
coverage_stats,
load_ticker_sector_map,
normalise_symbol,
save_ticker_sector_map,
sector_to_etf,
)
SP500_CSV_URL = (
"https://raw.githubusercontent.com/datasets/s-and-p-500-companies/"
"master/data/constituents.csv"
)
FMP_STABLE = "https://financialmodelingprep.com/stable"
def _parse_args() -> argparse.Namespace:
p = argparse.ArgumentParser(description=__doc__)
p.add_argument(
"--snapshot",
default="backtest_snapshots/prod.sqlite",
help="Snapshot whose tickers define the universe.",
)
p.add_argument(
"--out",
default=str(DEFAULT_SECTOR_MAP_PATH),
help="Output JSON path.",
)
p.add_argument(
"--fmp-limit",
type=int,
default=200,
help="Max FMP profile requests this run (free-tier cushion).",
)
p.add_argument(
"--skip-fmp",
action="store_true",
help="Only use public SP500 CSV + existing map.",
)
p.add_argument("--sleep", type=float, default=0.35, help="Pause between FMP calls.")
return p.parse_args()
def _snapshot_symbols(snapshot: Path) -> list[str]:
engine = create_engine(f"sqlite:///{snapshot.resolve().as_posix()}", future=True)
try:
with engine.connect() as conn:
rows = conn.execute(text("SELECT symbol FROM tickers ORDER BY symbol")).fetchall()
finally:
engine.dispose()
return [normalise_symbol(r[0]) for r in rows if r[0]]
def _fetch_sp500_map() -> dict[str, str]:
with httpx.Client(timeout=60.0, follow_redirects=True) as client:
resp = client.get(SP500_CSV_URL)
resp.raise_for_status()
reader = csv.DictReader(io.StringIO(resp.text))
out: dict[str, str] = {}
for row in reader:
sym = normalise_symbol(row.get("Symbol") or "")
sector = (row.get("GICS Sector") or "").strip()
if sym and sector:
out[sym] = sector
return out
async def _fmp_profile_sector(client: httpx.AsyncClient, api_key: str, symbol: str) -> str | None:
resp = await client.get(
f"{FMP_STABLE}/profile",
params={"symbol": symbol, "apikey": api_key},
)
if resp.status_code == 429:
raise RuntimeError(f"FMP rate limited on {symbol}")
if resp.status_code == 402:
return None
resp.raise_for_status()
data = resp.json()
if isinstance(data, list):
data = data[0] if data else {}
if not isinstance(data, dict):
return None
sector = (data.get("sector") or data.get("industry") or "").strip()
# industry alone is not a GICS sector — only accept if we can map to an ETF
if sector and sector_to_etf(sector):
return sector
# FMP sometimes returns industry under sector when sector missing; try sector field only
sec = (data.get("sector") or "").strip()
return sec or None
async def _fill_from_fmp(
missing: list[str],
*,
api_key: str,
limit: int,
sleep_s: float,
) -> tuple[dict[str, str], int]:
filled: dict[str, str] = {}
used = 0
async with httpx.AsyncClient(timeout=30.0) as client:
for sym in missing:
if used >= limit:
break
try:
sector = await _fmp_profile_sector(client, api_key, sym)
except Exception as exc:
print(f" FMP fail {sym}: {exc}")
used += 1
await asyncio.sleep(sleep_s)
continue
used += 1
if sector:
filled[sym] = sector
print(f" FMP {sym}{sector}")
else:
print(f" FMP {sym} → (no sector)")
if sleep_s > 0:
await asyncio.sleep(sleep_s)
return filled, used
async def _main() -> None:
args = _parse_args()
snapshot = Path(args.snapshot)
if not snapshot.exists():
raise SystemExit(f"Snapshot not found: {snapshot}")
symbols = _snapshot_symbols(snapshot)
print(f"Universe: {len(symbols)} symbols from {snapshot}")
existing = load_ticker_sector_map(args.out)
print(f"Existing map entries: {len(existing)}")
print("Fetching public S&P 500 sector CSV…")
sp500 = _fetch_sp500_map()
print(f" SP500 CSV rows: {len(sp500)}")
mapping = dict(existing)
from_sp500 = 0
for sym in symbols:
if sym in mapping:
continue
if sym in sp500:
mapping[sym] = sp500[sym]
from_sp500 += 1
print(f" Newly filled from SP500 CSV: {from_sp500}")
missing = [s for s in symbols if s not in mapping]
fmp_used = 0
from_fmp = 0
if missing and not args.skip_fmp:
from app.config import settings
if not settings.fmp_api_key:
print("WARNING: FMP key missing; leaving gaps unfilled")
else:
print(f"FMP fill for {len(missing)} missing (limit={args.fmp_limit})…")
filled, fmp_used = await _fill_from_fmp(
missing,
api_key=settings.fmp_api_key,
limit=int(args.fmp_limit),
sleep_s=float(args.sleep),
)
mapping.update(filled)
from_fmp = len(filled)
still_missing = [s for s in symbols if s not in mapping]
stats = coverage_stats(symbols, mapping)
meta = {
"built_at": datetime.now(timezone.utc).isoformat(),
"snapshot": str(snapshot.resolve()),
"from_existing": len(existing),
"from_sp500_csv": from_sp500,
"from_fmp": from_fmp,
"fmp_requests": fmp_used,
"still_missing": still_missing,
"coverage": {
k: stats[k]
for k in ("universe", "mapped", "mapped_pct", "with_etf", "by_sector")
},
}
out_path = save_ticker_sector_map(mapping, args.out, meta=meta)
print(f"Wrote {out_path}")
print(json.dumps(meta["coverage"], indent=2))
if still_missing:
print(f"Still missing ({len(still_missing)}): {still_missing[:40]}")
if len(still_missing) > 40:
print(f" … +{len(still_missing) - 40} more")
if __name__ == "__main__":
asyncio.run(_main())
-187
View File
@@ -1,187 +0,0 @@
"""Fetch the 11 SPDR sector ETFs into a snapshot's ``benchmark_prices``.
Research-only. Sector ETFs are auxiliary series (like SPY) — they must not
enter the tradable ticker universe or candidate replay. Storing them in
``benchmark_prices`` keeps that invariant.
Also refreshes SPY on the same window so residual factors share a calendar.
Example
-------
python scripts/fetch_sector_etfs_to_snapshot.py \\
--snapshot backtest_snapshots/prod.sqlite --history-days 2200
"""
from __future__ import annotations
import argparse
import asyncio
import sys
import time
from datetime import date, timedelta
from pathlib import Path
from sqlalchemy import create_engine, text
ROOT = Path(__file__).resolve().parents[1]
if str(ROOT) not in sys.path:
sys.path.insert(0, str(ROOT))
from app.ssl_bootstrap import bootstrap_ssl # noqa: E402
bootstrap_ssl()
from app.services.sector_map import SECTOR_ETFS # noqa: E402
def _parse_args() -> argparse.Namespace:
p = argparse.ArgumentParser(description=__doc__)
p.add_argument("--snapshot", default="backtest_snapshots/prod.sqlite")
p.add_argument(
"--history-days",
type=int,
default=2200,
help="Lookback calendar days (default ~6y; covers 5y snapshot + cushion).",
)
p.add_argument("--sleep", type=float, default=0.25)
p.add_argument(
"--symbols",
default=None,
help="Comma-separated override (default: SPY + 11 sector ETFs).",
)
return p.parse_args()
async def _fetch_and_upsert(
engine,
provider,
symbol: str,
start: date,
end: date,
*,
sleep_s: float,
) -> int:
from app.exceptions import ProviderError, RateLimitError
for attempt in range(5):
try:
bars = await provider.fetch_ohlcv(symbol, start, end)
break
except RateLimitError:
wait = min(60.0, 2.0 ** attempt)
print(f" rate limited {symbol}; sleep {wait:.0f}s")
await asyncio.sleep(wait)
bars = []
except ProviderError as exc:
if attempt + 1 >= 5:
raise
await asyncio.sleep(1.0)
print(f" retry {symbol}: {exc}")
bars = []
else:
bars = []
if sleep_s > 0:
await asyncio.sleep(sleep_s)
if not bars:
print(f" {symbol}: empty")
return 0
written = 0
with engine.begin() as conn:
for bar in bars:
d = bar.date.isoformat() if hasattr(bar.date, "isoformat") else str(bar.date)
close = float(bar.close)
existing = conn.execute(
text(
"SELECT id, close FROM benchmark_prices "
"WHERE symbol = :sym AND date = :d"
),
{"sym": symbol, "d": d},
).fetchone()
if existing is None:
# id is INTEGER PK — let sqlite autoincrement if possible
conn.execute(
text(
"INSERT INTO benchmark_prices (symbol, date, close) "
"VALUES (:sym, :d, :c)"
),
{"sym": symbol, "d": d, "c": close},
)
written += 1
elif abs(float(existing[1]) - close) > 1e-9:
conn.execute(
text(
"UPDATE benchmark_prices SET close = :c WHERE id = :id"
),
{"c": close, "id": int(existing[0])},
)
written += 1
print(f" {symbol}: {len(bars)} bars, {written} rows written/updated")
return written
async def _main() -> None:
args = _parse_args()
snapshot = Path(args.snapshot)
if not snapshot.exists():
raise SystemExit(f"Snapshot not found: {snapshot}")
from app.config import settings
from app.providers.alpaca import AlpacaOHLCVProvider
if not settings.alpaca_api_key or not settings.alpaca_api_secret:
raise SystemExit("ALPACA_API_KEY / ALPACA_API_SECRET required")
if args.symbols:
symbols = [s.strip().upper() for s in args.symbols.split(",") if s.strip()]
else:
symbols = ["SPY", *SECTOR_ETFS]
end = date.today()
start = end - timedelta(days=int(args.history_days))
provider = AlpacaOHLCVProvider(settings.alpaca_api_key, settings.alpaca_api_secret)
engine = create_engine(
f"sqlite:///{snapshot.resolve().as_posix()}",
future=True,
)
print(f"Snapshot: {snapshot}")
print(f"Window: {start}{end}")
print(f"Symbols: {symbols}")
t0 = time.monotonic()
total = 0
try:
for sym in symbols:
n = await _fetch_and_upsert(
engine, provider, sym, start, end, sleep_s=float(args.sleep)
)
total += n
finally:
engine.dispose()
# Summary counts
engine = create_engine(
f"sqlite:///{snapshot.resolve().as_posix()}",
future=True,
)
try:
with engine.connect() as conn:
rows = conn.execute(
text(
"SELECT symbol, COUNT(*), MIN(date), MAX(date) "
"FROM benchmark_prices GROUP BY symbol ORDER BY symbol"
)
).fetchall()
finally:
engine.dispose()
print(f"Done in {(time.monotonic() - t0) / 60:.1f}m; rows touched={total}")
for sym, n, d0, d1 in rows:
print(f" {sym}: n={n} {d0}{d1}")
if __name__ == "__main__":
asyncio.run(_main())
-26
View File
@@ -466,11 +466,6 @@ async def _run_2b_ic(
os.environ["BACKTEST_SNAPSHOT_OFFLINE"] = "1"
os.environ["BACKTEST_SIGNAL_EVAL_ONLY"] = "1"
# Load sector map if present so sector signals also appear (side-by-side optional).
if Path("data/research/ticker_sector_map.json").exists():
os.environ["BACKTEST_SECTOR_MAP_PATH"] = str(
Path("data/research/ticker_sector_map.json").resolve()
)
settings.backtest_workers = workers
engine = create_async_engine(_sqlite_url(snapshot), pool_pre_ping=True)
@@ -484,18 +479,6 @@ async def _run_2b_ic(
(await db.execute(select(Ticker).order_by(Ticker.symbol))).scalars()
)
spy = await load_benchmark_closes(db, "SPY")
sector_etf: dict[str, dict] = {}
try:
from app.services.sector_map import SECTOR_ETFS, load_ticker_sector_map
symbol_to_sector = load_ticker_sector_map()
for etf in SECTOR_ETFS:
series = await load_benchmark_closes(db, etf)
if series:
sector_etf[etf] = series
except Exception:
symbol_to_sector = {}
sector_etf = {}
prices: dict[str, tuple] = {}
for idx, t in enumerate(tickers):
@@ -521,9 +504,6 @@ async def _run_2b_ic(
],
spy,
symbol=t.symbol,
sector_etf_closes=bt._sector_etf_closes_for_symbol(
t.symbol, symbol_to_sector, sector_etf
),
)
for name, weeks in series.items():
for wk, pairs in weeks.items():
@@ -533,9 +513,6 @@ async def _run_2b_ic(
if not quiet:
print()
if symbol_to_sector:
bt._inject_sector_demeaned_momentum(collected, symbol_to_sector)
# SUE series.
events_by_sym: dict[str, list[dict]] = defaultdict(list)
for ev in events:
@@ -709,8 +686,6 @@ async def _run_2b_ic(
for name in (
"mom_12_1",
"mom_12_1_resid",
"mom_12_1_sector_resid",
"mom_12_1_sector_demeaned",
"sue_latest",
"fip_id",
)
@@ -784,7 +759,6 @@ def _write_md(path: Path, payload: dict) -> None:
"mom_12_1",
"mom_12_1_resid",
"sue_latest",
"mom_12_1_sector_resid",
"fip_id",
):
r = side.get(name) or {}
-480
View File
@@ -1,480 +0,0 @@
"""History-depth extension research (local / MacBook).
Phases
------
coverage — bars per calendar year; no rebuild
harness — race-guard snapshot, full signal_eval, era split pre/post-2021
Does not retune production knobs. Does not modify scheduler/gates.
Example
-------
python scripts/run_history_depth_research.py --phase coverage \\
--snapshot backtest_snapshots/prod.sqlite
python scripts/run_history_depth_research.py --phase harness \\
--snapshot backtest_snapshots/research.sqlite --workers 8 --allow-spawn
"""
from __future__ import annotations
import argparse
import asyncio
import json
import os
import sys
from collections import defaultdict
from datetime import date, datetime
from pathlib import Path
from typing import Any
from sqlalchemy import create_engine, text
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 app.ssl_bootstrap import bootstrap_ssl # noqa: E402
bootstrap_ssl()
ERA_SPLIT = date(2021, 1, 1)
SURVIVORSHIP_BANNER = (
"SURVIVORSHIP BIAS: today's constituents backfilled historically. "
"Absolute Sharpe/CAGR levels on deep history are optimistic. "
"Use RELATIVE signal IC comparisons and era stability only — not levels."
)
def _sqlite_url(path: Path) -> str:
return f"sqlite+aiosqlite:///{path.resolve().as_posix()}"
def _parse_args() -> argparse.Namespace:
p = argparse.ArgumentParser(description=__doc__)
p.add_argument("--phase", choices=("coverage", "harness", "all"), default="all")
p.add_argument("--snapshot", default="backtest_snapshots/research.sqlite")
p.add_argument("--workers", type=int, default=8)
p.add_argument("--allow-spawn", action="store_true")
p.add_argument("--quiet", action="store_true")
p.add_argument("--out", default=None)
return p.parse_args()
def _coverage_report(snapshot: Path) -> dict[str, Any]:
engine = create_engine(
f"sqlite:///{snapshot.resolve().as_posix()}",
future=True,
)
try:
with engine.connect() as conn:
ticker_n = int(conn.execute(text("SELECT COUNT(*) FROM tickers")).scalar_one())
ohlcv_n = int(
conn.execute(text("SELECT COUNT(*) FROM ohlcv_records")).scalar_one()
)
d_range = conn.execute(
text("SELECT MIN(date), MAX(date) FROM ohlcv_records")
).fetchone()
# Bars per calendar year (global).
by_year = conn.execute(
text(
"""
SELECT substr(date, 1, 4) AS y, COUNT(*) AS n,
COUNT(DISTINCT ticker_id) AS tickers
FROM ohlcv_records
GROUP BY substr(date, 1, 4)
ORDER BY y
"""
)
).fetchall()
# Per-symbol min/max date + bar count (summary percentiles).
per_sym = conn.execute(
text(
"""
SELECT t.symbol, COUNT(*) AS n, MIN(o.date), MAX(o.date)
FROM ohlcv_records o
JOIN tickers t ON t.id = o.ticker_id
GROUP BY t.symbol
"""
)
).fetchall()
finally:
engine.dispose()
ns = sorted(int(r[1]) for r in per_sym)
def pct(p: float) -> int | None:
if not ns:
return None
i = int(round(p * (len(ns) - 1)))
return ns[i]
starts = sorted(str(r[2]) for r in per_sym if r[2])
start_hist: dict[str, int] = defaultdict(int)
for s in starts:
start_hist[s[:4]] += 1
return {
"snapshot": str(snapshot.resolve()),
"ticker_count": ticker_n,
"ohlcv_row_count": ohlcv_n,
"date_range": {"min": d_range[0], "max": d_range[1]},
"bars_per_year": [
{"year": y, "bars": n, "tickers_with_bars": t} for y, n, t in by_year
],
"bars_per_symbol": {
"min": ns[0] if ns else None,
"p10": pct(0.10),
"p50": pct(0.50),
"p90": pct(0.90),
"max": ns[-1] if ns else None,
},
"symbols_by_start_year": dict(sorted(start_hist.items())),
"note": (
"Where ticker counts drop in early years, the feed (or listing history) "
"thins — do not treat those years as a full 505-name cross-section."
),
"survivorship_banner": SURVIVORSHIP_BANNER,
}
def _assert_complete(snapshot: Path) -> dict[str, Any]:
from scripts.research_snapshot_manifest import ( # type: ignore
assert_research_snapshot_complete,
load_manifest,
)
m = load_manifest(snapshot)
if m is None:
# Prod snapshot may lack manifest; still require healthy bar depth.
eng = create_engine(
f"sqlite:///{snapshot.resolve().as_posix()}",
future=True,
)
try:
with eng.connect() as conn:
avg = conn.execute(
text(
"""
SELECT AVG(c) FROM (
SELECT COUNT(*) AS c FROM ohlcv_records GROUP BY ticker_id
)
"""
)
).scalar_one()
finally:
eng.dispose()
if avg is None or float(avg) < 400:
raise SystemExit(
f"No completion manifest and avg bars={avg} look short. "
"Rebuild research.sqlite via extend_snapshot_universe.py"
)
return {"manifest": None, "avg_bars": float(avg), "ok": True}
return {"manifest": assert_research_snapshot_complete(snapshot), "ok": True}
async def _harness(snapshot: Path, *, workers: int, quiet: bool) -> dict[str, Any]:
from app.config import settings
from app.services.backtest_service import run_backtest
os.environ["BACKTEST_SNAPSHOT_OFFLINE"] = "1"
os.environ["BACKTEST_SIGNAL_EVAL_ONLY"] = "1"
if Path("data/research/ticker_sector_map.json").exists():
os.environ["BACKTEST_SECTOR_MAP_PATH"] = str(
Path("data/research/ticker_sector_map.json").resolve()
)
settings.backtest_workers = workers
engine = create_async_engine(_sqlite_url(snapshot), pool_pre_ping=True)
Session = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
def progress(done: int, total: int, symbol: str) -> None:
if quiet:
return
print(f" progress {done}/{total} {symbol}", end="\r", flush=True)
try:
async with Session() as db:
report = await run_backtest(db, progress_cb=progress, cadence="weekly")
finally:
await engine.dispose()
if not quiet:
print()
signal_eval = report.get("signal_eval") or []
# Era-split IC: recompute from collected is not available post-run.
# Approximate via second pass is expensive; instead document that era split
# requires collecting weekly ICs. We re-run evaluation if the report embeds
# nothing — for v1, call internal collection is too heavy to duplicate.
# Lightweight approach: mark era_split as requiring BACKTEST with custom
# filter — implemented below by re-scoring from a dedicated collection pass.
era = await _era_split_ics(snapshot, workers=workers, quiet=quiet)
return {
"survivorship_banner": SURVIVORSHIP_BANNER,
"signal_eval": signal_eval,
"era_split": era,
"params": report.get("params"),
"tickers": report.get("tickers"),
"generated_at_run": report.get("generated_at"),
}
async def _era_split_ics(
snapshot: Path, *, workers: int, quiet: bool
) -> dict[str, Any]:
"""Collect weekly signal series and evaluate pre/post ERA_SPLIT separately."""
from app.config import settings
from app.services import backtest_service as bt
from app.services.benchmark_service import load_benchmark_closes
from app.models.ticker import Ticker
from sqlalchemy import select
from collections import defaultdict as dd
os.environ["BACKTEST_SNAPSHOT_OFFLINE"] = "1"
settings.backtest_workers = max(1, workers)
engine = create_async_engine(_sqlite_url(snapshot), pool_pre_ping=True)
Session = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
collected: dict = dd(lambda: dd(list))
try:
async with Session() as db:
tickers = list(
(await db.execute(select(Ticker).order_by(Ticker.symbol))).scalars()
)
spy = await load_benchmark_closes(db, "SPY")
symbol_to_sector = {}
sector_etf: dict = {}
try:
from app.services.sector_map import (
SECTOR_ETFS,
load_ticker_sector_map,
)
symbol_to_sector = load_ticker_sector_map()
for etf in SECTOR_ETFS:
series = await load_benchmark_closes(db, etf)
if series:
sector_etf[etf] = series
except Exception:
pass
for idx, t in enumerate(tickers):
if not quiet and idx % 100 == 0:
print(f" era-collect {idx}/{len(tickers)}", end="\r", flush=True)
cols = await bt._fetch_columns(db, t.symbol)
if cols is None:
continue
records = [
type(
"R",
(),
{
"date": date.fromordinal(int(cols[0][i])),
"close": cols[4][i],
"high": cols[2][i],
"volume": cols[5][i] if len(cols) > 5 else 0,
},
)()
for i in range(len(cols[0]))
]
series = bt._signal_series(
records,
spy,
symbol=t.symbol,
sector_etf_closes=bt._sector_etf_closes_for_symbol(
t.symbol, symbol_to_sector, sector_etf
),
)
for name, weeks in series.items():
for wk, pairs in weeks.items():
collected[name][wk].extend(pairs)
if symbol_to_sector:
bt._inject_sector_demeaned_momentum(collected, symbol_to_sector)
finally:
await engine.dispose()
if not quiet:
print()
def _filter_era(coll: dict, *, pre: bool) -> dict:
out: dict = dd(lambda: dd(list))
for name, weeks in coll.items():
for wk, recs in weeks.items():
# ISO week key (year, week) — approximate era by ISO year.
year = int(wk[0]) if isinstance(wk, tuple) else int(str(wk)[:4])
if pre and year >= ERA_SPLIT.year:
continue
if not pre and year < ERA_SPLIT.year:
continue
out[name][wk].extend(recs)
return out
pre_eval = bt._signal_evaluation(_filter_era(collected, pre=True))
post_eval = bt._signal_evaluation(_filter_era(collected, pre=False))
full_eval = bt._signal_evaluation(collected)
def _index(rows: list[dict]) -> dict[str, dict]:
return {r["signal"]: r for r in rows}
return {
"era_split_date": ERA_SPLIT.isoformat(),
"note": "Diagnostic only — not a tuning input. Nested lookbacks are not OOS.",
"full": _index(full_eval),
"pre_2021": _index(pre_eval),
"post_2021": _index(post_eval),
}
def _write_md(path: Path, payload: dict) -> None:
pre = path.read_text(encoding="utf-8") if path.exists() else ""
marker = "## Results"
idx = pre.find(marker)
header = pre[:idx] if idx >= 0 else pre.split("## Verdict")[0]
lines = [
header.rstrip(),
"",
"## Results",
"",
f"Generated: `{payload.get('generated_at')}`",
"",
f"> **{SURVIVORSHIP_BANNER}**",
"",
"### Coverage",
"",
f"```json\n{json.dumps(payload.get('coverage') or {}, indent=2, default=str)}\n```",
"",
"### Race guard",
"",
f"```json\n{json.dumps(payload.get('race_guard') or {}, indent=2, default=str)}\n```",
"",
"### Signal IC (full extended window)",
"",
]
harness = payload.get("harness") or {}
rows = harness.get("signal_eval") or []
if rows:
lines.extend([
"| signal | mean_ic | ic_t_stat | weeks | avg_N | reliable |",
"|---|---:|---:|---:|---:|---|",
])
for r in rows:
lines.append(
f"| {r.get('signal')} | {r.get('mean_ic')} | {r.get('ic_t_stat')} | "
f"{r.get('weeks')} | {r.get('avg_cross_section')} | {r.get('reliable')} |"
)
else:
lines.append("_Harness not run this pass._")
era = (harness.get("era_split") or {})
lines.extend(["", "### Era split (diagnostic only)", ""])
if era:
for label in ("full", "pre_2021", "post_2021"):
block = era.get(label) or {}
lines.append(f"#### {label}")
lines.append("")
lines.append("| signal | mean_ic | t | weeks | N |")
lines.append("|---|---:|---:|---:|---:|")
for name in sorted(block):
r = block[name]
lines.append(
f"| {name} | {r.get('mean_ic')} | {r.get('ic_t_stat')} | "
f"{r.get('weeks')} | {r.get('avg_cross_section')} |"
)
lines.append("")
else:
lines.append("_No era split._")
lines.extend([
"",
"## Verdict",
"",
f"**{payload.get('verdict')}**",
"",
payload.get("verdict_detail") or "",
"",
"## What a human must decide next",
"",
payload.get("human_next")
or "- Do not retune production knobs from this report without review.",
"",
f"Artifacts: `{payload.get('report_path')}`",
"",
])
path.write_text("\n".join(lines) + "\n", encoding="utf-8")
async def _main() -> None:
args = _parse_args()
snapshot = Path(args.snapshot)
if not snapshot.exists():
raise SystemExit(f"Missing snapshot: {snapshot}")
if args.allow_spawn:
os.environ["BACKTEST_ALLOW_SPAWN"] = "1"
coverage = None
race = None
harness = None
if args.phase in ("coverage", "all"):
print("Coverage probe…")
coverage = _coverage_report(snapshot)
print(
f" tickers={coverage['ticker_count']} ohlcv={coverage['ohlcv_row_count']} "
f"range={coverage['date_range']}"
)
for row in coverage["bars_per_year"]:
print(
f" year {row['year']}: bars={row['bars']} "
f"tickers={row['tickers_with_bars']}"
)
if args.phase in ("harness", "all"):
print("Race guard…")
race = _assert_complete(snapshot)
print(f" ok={race.get('ok')}")
print("Full harness + era split (LONG)…")
print(f" {SURVIVORSHIP_BANNER}")
harness = await _harness(
snapshot, workers=args.workers, quiet=args.quiet
)
stamp = datetime.now().strftime("%Y%m%d-%H%M%S")
out = (
Path(args.out)
if args.out
else Path("reports") / f"history-depth-{stamp}.json"
)
payload = {
"generated_at": datetime.now().isoformat(),
"survivorship_banner": SURVIVORSHIP_BANNER,
"coverage": coverage,
"race_guard": race,
"harness": harness,
"verdict": "PENDING_HUMAN" if harness else "COVERAGE_ONLY",
"verdict_detail": (
"Harness complete — human interprets relative IC / era stability. "
"No production retune from this artifact."
if harness
else "Coverage probe only; run --phase harness after deep rebuild."
),
"human_next": (
"- Compare sector residual vs market residual across eras.\n"
"- If pre-2021 IC collapses, park Task 1 wire-in.\n"
"- Do not retune production knobs on deep history levels."
),
"report_path": str(out.as_posix()),
}
out.parent.mkdir(parents=True, exist_ok=True)
out.write_text(json.dumps(payload, indent=2, default=str) + "\n", encoding="utf-8")
md = Path("docs/research/history-depth-extension.md")
_write_md(md, payload)
out.with_suffix(".md").write_text(md.read_text(encoding="utf-8"), encoding="utf-8")
print(f"Wrote {out}")
print(f"Wrote {md}")
if __name__ == "__main__":
asyncio.run(_main())
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
+46 -198
View File
@@ -1,99 +1,70 @@
#!/usr/bin/env bash
# Tier-1 alpha research runner for a high-CPU MacBook (local only).
# Research helpers for a high-CPU MacBook (local only).
#
# Prerequisites
# - git checkout research/earnings-gap-and-sue (or later research branch)
# - .env with ALPACA_* (required for OHLCV/ETFs); FMP_* for earnings resume;
# optional ALPHA_VANTAGE_* as earnings fallback
# - Python venv with project deps installed
# - backtest_snapshots/prod.sqlite present (gitignored — copy or rebuild)
# Kept after Tier-1 cleanup:
# --ssl-check diagnose corporate CA / proxy
# --earnings-only resume FMP earnings backfill + 2a/2b (parked)
# --prod-book-matrix re-run 505 vs liquid universe × horizon book matrix
#
# Prerequisites: git checkout research branch, .env, deep research.sqlite for
# book matrix, combined-ca-bundle.pem or certifi when on corp network.
#
# Usage
# chmod +x scripts/run_tier1_macbook.sh
# ./scripts/run_tier1_macbook.sh # coverage + deep rebuild + harness
# ./scripts/run_tier1_macbook.sh --all # earnings resume + full depth pipeline
# ./scripts/run_tier1_macbook.sh --earnings-only # multi-day FMP backfill + re-run 2a/2b
# ./scripts/run_tier1_macbook.sh --harness-only # skip rebuild; race-guard + IC only
# ./scripts/run_tier1_macbook.sh --coverage-only # bars-per-year probe only
# ./scripts/run_tier1_macbook.sh --sector-resid-deep # deepen shallow + ONE masked grade
# ./scripts/run_tier1_macbook.sh --prod-book-matrix # 4-arm universe×horizon book matrix
#
# Does NOT touch production Postgres, scheduler, gates, or prod config.
# ./scripts/run_tier1_macbook.sh --ssl-check
# ./scripts/run_tier1_macbook.sh --prod-book-matrix
set -euo pipefail
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
cd "$ROOT"
# --- defaults (override via flags or env) ---
PROD_SNAP="${PROD_SNAP:-backtest_snapshots/prod.sqlite}"
RESEARCH_SNAP="${RESEARCH_SNAP:-backtest_snapshots/research.sqlite}"
HISTORY_DAYS="${HISTORY_DAYS:-5000}"
MIN_BARS="${MIN_BARS:-260}"
PROD_SNAP="${PROD_SNAP:-backtest_snapshots/prod.sqlite}"
WORKERS="${WORKERS:-8}"
ALPACA_SLEEP="${ALPACA_SLEEP:-0.15}"
FMP_LIMIT="${FMP_LIMIT:-250}"
FMP_SLEEP="${FMP_SLEEP:-0.35}"
PYTHON="${PYTHON:-python3}"
USE_CORP_PROXY="${USE_CORP_PROXY:-0}"
PHASE="depth" # depth | all | earnings | harness | coverage | ssl | sector-resid-deep | prod-book
PHASE=""
usage() {
sed -n '2,25p' "$0" | sed 's/^# \?//'
cat <<'EOF'
SSL / network (corporate MacBook)
SSL errors usually mean the corp root CA is missing from Python.
1) Put combined-ca-bundle.pem in the repo root OR $HOME
2) Or: export SSL_CERT_FILE=/path/to/combined-ca-bundle.pem
3) Behind corp proxy: USE_CORP_PROXY=1 ./scripts/run_tier1_macbook.sh
4) Diagnose: ./scripts/run_tier1_macbook.sh --ssl-check
EOF
sed -n '2,16p' "$0" | sed 's/^# \?//'
exit "${1:-0}"
}
while [[ $# -gt 0 ]]; do
case "$1" in
--all) PHASE=all; shift ;;
--earnings-only) PHASE=earnings; shift ;;
--harness-only) PHASE=harness; shift ;;
--coverage-only) PHASE=coverage; shift ;;
--depth) PHASE=depth; shift ;;
--ssl-check) PHASE=ssl; shift ;;
--sector-resid-deep) PHASE=sector_resid_deep; shift ;;
--earnings-only) PHASE=earnings; shift ;;
--prod-book-matrix) PHASE=prod_book; shift ;;
--corp-proxy) USE_CORP_PROXY=1; shift ;;
--prod-snap) PROD_SNAP="$2"; shift 2 ;;
--research-snap) RESEARCH_SNAP="$2"; shift 2 ;;
--history-days) HISTORY_DAYS="$2"; shift 2 ;;
--workers) WORKERS="$2"; shift 2 ;;
--fmp-limit) FMP_LIMIT="$2"; shift 2 ;;
--python) PYTHON="$2"; shift 2 ;;
-h|--help) usage 0 ;;
*) echo "Unknown flag: $1" >&2; usage 1 ;;
esac
done
if [[ -z "$PHASE" ]]; then
echo "Pick a phase: --ssl-check | --earnings-only | --prod-book-matrix" >&2
usage 1
fi
if [[ -x .venv/bin/python ]]; then
PYTHON=".venv/bin/python"
elif command -v "$PYTHON" >/dev/null 2>&1; then
:
else
echo "ERROR: no Python found (tried .venv/bin/python and $PYTHON)" >&2
echo "ERROR: no Python found" >&2
exit 1
fi
log() { printf '\n==> %s\n' "$*"; }
die() { echo "ERROR: $*" >&2; exit 1; }
need_file() { [[ -f "$1" ]] || die "missing $1"; }
# ---------------------------------------------------------------------------
# TLS bootstrap — same corp CA path the FastAPI app uses
# ---------------------------------------------------------------------------
setup_ssl() {
export USE_CORP_PROXY
# Prefer explicit env, then repo / home corporate bundle, then certifi.
if [[ -z "${SSL_CERT_FILE:-}" ]]; then
if [[ -f "$ROOT/combined-ca-bundle.pem" ]]; then
export SSL_CERT_FILE="$ROOT/combined-ca-bundle.pem"
@@ -101,191 +72,68 @@ setup_ssl() {
export SSL_CERT_FILE="$HOME/combined-ca-bundle.pem"
fi
fi
if [[ -n "${SSL_CERT_FILE:-}" && -f "$SSL_CERT_FILE" ]]; then
export REQUESTS_CA_BUNDLE="$SSL_CERT_FILE"
export CURL_CA_BUNDLE="$SSL_CERT_FILE"
export REQUESTS_CA_BUNDLE="$SSL_CERT_FILE" CURL_CA_BUNDLE="$SSL_CERT_FILE"
log "SSL CA bundle: $SSL_CERT_FILE"
else
# Fall back to certifi if installed
local certifi_path
certifi_path="$("$PYTHON" -c 'import certifi; print(certifi.where())' 2>/dev/null || true)"
if [[ -n "$certifi_path" && -f "$certifi_path" ]]; then
export SSL_CERT_FILE="$certifi_path"
export REQUESTS_CA_BUNDLE="$certifi_path"
export CURL_CA_BUNDLE="$certifi_path"
export SSL_CERT_FILE="$certifi_path" REQUESTS_CA_BUNDLE="$certifi_path" CURL_CA_BUNDLE="$certifi_path"
log "SSL CA bundle (certifi): $SSL_CERT_FILE"
else
log "WARNING: no CA bundle found — SSL may fail on corp networks"
log " Copy combined-ca-bundle.pem to $ROOT/ or \$HOME/"
log " Or: export SSL_CERT_FILE=/path/to/combined-ca-bundle.pem"
fi
fi
if [[ "$USE_CORP_PROXY" == "1" ]]; then
export HTTP_PROXY="${HTTP_PROXY:-http://aproxy.corproot.net:8080}"
export HTTPS_PROXY="${HTTPS_PROXY:-http://aproxy.corproot.net:8080}"
export NO_PROXY="${NO_PROXY:-corproot.net,sharedtcs.net,127.0.0.1,localhost,bix.swisscom.com,swisscom.com}"
export NO_PROXY="${NO_PROXY:-corproot.net,sharedtcs.net,127.0.0.1,localhost}"
export http_proxy="$HTTP_PROXY" https_proxy="$HTTPS_PROXY" no_proxy="$NO_PROXY"
log "Corp proxy enabled: $HTTPS_PROXY"
fi
# Ensure Python process sees the same bootstrap (patches ssl for alpaca-py).
export PYTHONPATH="${ROOT}${PYTHONPATH:+:$PYTHONPATH}"
}
ssl_check() {
setup_ssl
log "SSL diagnostic"
"$PYTHON" - <<'PY'
from app.ssl_bootstrap import bootstrap_ssl, ssl_status
import json
import urllib.request
ca = bootstrap_ssl()
import json, urllib.request
print(json.dumps(ssl_status(), indent=2))
print("bootstrap_ssl ->", ca)
urls = [
print("bootstrap ->", bootstrap_ssl())
for url in (
"https://data.alpaca.markets/v2/stocks/SPY/bars?timeframe=1Day&limit=1",
"https://financialmodelingprep.com/stable/profile?symbol=AAPL",
"https://www.alphavantage.co/query?function=TIME_SERIES_DAILY&symbol=IBM",
]
for url in urls:
):
try:
req = urllib.request.Request(url, headers={"User-Agent": "signal-platform-ssl-check"})
req = urllib.request.Request(url, headers={"User-Agent": "ssl-check"})
with urllib.request.urlopen(req, timeout=20) as resp:
print(f"OK {resp.status} {url[:60]}...")
print(f"OK {resp.status} {url[:60]}")
except Exception as exc:
print(f"FAIL {type(exc).__name__}: {exc}")
print(f" {url[:80]}")
PY
}
need_file() {
[[ -f "$1" ]] || die "missing $1"
}
require_prod() {
need_file "$PROD_SNAP"
}
run_earnings() {
require_prod
log "Earnings backfill (FMP free tier ~${FMP_LIMIT}/day; resume-safe)"
"$PYTHON" scripts/backfill_earnings_events.py \
--snapshot "$PROD_SNAP" \
--provider fmp \
--force-symbol \
--limit "$FMP_LIMIT" \
--sleep "$FMP_SLEEP"
log "Earnings research 2a+2b (report only; no filters shipped)"
"$PYTHON" scripts/run_earnings_research.py \
--snapshot "$PROD_SNAP" \
--workers "$WORKERS" \
--allow-spawn
}
run_coverage() {
require_prod
log "Coverage probe (bars per year) on $PROD_SNAP"
"$PYTHON" scripts/run_history_depth_research.py \
--phase coverage \
--snapshot "$PROD_SNAP"
}
run_rebuild() {
require_prod
log "Deep rebuild $PROD_SNAP$RESEARCH_SNAP (history-days=$HISTORY_DAYS)"
log "SURVIVORSHIP: today's constituents backfilled — relative IC only, not levels"
"$PYTHON" scripts/extend_snapshot_universe.py \
--source "$PROD_SNAP" \
--output "$RESEARCH_SNAP" \
--force-copy \
--history-days "$HISTORY_DAYS" \
--min-bars "$MIN_BARS" \
--sleep "$ALPACA_SLEEP"
log "Refresh SPY + 11 sector ETFs on research snapshot"
"$PYTHON" scripts/fetch_sector_etfs_to_snapshot.py \
--snapshot "$RESEARCH_SNAP" \
--history-days "$HISTORY_DAYS"
log "Also deepen sector ETFs on prod snapshot (for local A/B parity)"
"$PYTHON" scripts/fetch_sector_etfs_to_snapshot.py \
--snapshot "$PROD_SNAP" \
--history-days "$HISTORY_DAYS"
}
run_harness() {
need_file "$RESEARCH_SNAP"
log "Race-guard + full signal harness + era split on $RESEARCH_SNAP"
"$PYTHON" scripts/run_history_depth_research.py \
--phase harness \
--snapshot "$RESEARCH_SNAP" \
--workers "$WORKERS" \
--allow-spawn
}
run_sector_resid_deep() {
need_file "$RESEARCH_SNAP"
log "Sector-resid deep test: deepen shallow symbols + ONE liquid-1500 masked grade"
log "Pre-registered PASS/FAIL only — thread ends after this run"
"$PYTHON" scripts/run_sector_resid_deep_test.py \
--snapshot "$RESEARCH_SNAP" \
--history-days "$HISTORY_DAYS" \
--sleep "$ALPACA_SLEEP" \
--workers "$WORKERS" \
--allow-spawn
}
run_prod_book_matrix() {
need_file "$RESEARCH_SNAP"
log "Production book × universe × horizon (4 arms, strategy unchanged)"
"$PYTHON" scripts/run_prod_book_universe_matrix.py \
--snapshot "$RESEARCH_SNAP" \
--workers "$WORKERS" \
--allow-spawn \
--candidate-cache reports/.cache/prod-book-universe-cands.pkl
}
log "cwd=$ROOT python=$PYTHON phase=$PHASE workers=$WORKERS"
setup_ssl
case "$PHASE" in
ssl)
ssl_check
;;
sector_resid_deep)
run_sector_resid_deep
ssl) ssl_check ;;
earnings)
need_file "$PROD_SNAP"
log "Earnings backfill + research (parked experiment)"
"$PYTHON" scripts/backfill_earnings_events.py \
--snapshot "$PROD_SNAP" --provider fmp --force-symbol \
--limit "$FMP_LIMIT" --sleep "$FMP_SLEEP"
"$PYTHON" scripts/run_earnings_research.py \
--snapshot "$PROD_SNAP" --workers "$WORKERS" --allow-spawn
;;
prod_book)
run_prod_book_matrix
;;
coverage)
run_coverage
;;
earnings)
run_earnings
;;
harness)
run_harness
;;
depth)
run_coverage
run_rebuild
run_harness
;;
all)
run_earnings
run_coverage
run_rebuild
run_harness
;;
*)
die "unknown phase $PHASE"
need_file "$RESEARCH_SNAP"
log "Production book universe × horizon matrix"
"$PYTHON" scripts/run_prod_book_universe_matrix.py \
--snapshot "$RESEARCH_SNAP" --workers "$WORKERS" --allow-spawn \
--candidate-cache reports/.cache/prod-book-universe-cands.pkl
;;
*) die "unknown phase $PHASE" ;;
esac
log "Done. Check reports/ and docs/research/history-depth-extension.md"
log "Commit reports on this machine if they look good, or copy them back to Windows."
log "Done."