[verified] Task 6: durable backtest run store + strict /api/v1/backtest/run route
This commit is contained in:
@@ -150,6 +150,10 @@ def create_app(config: dict[str, Any] | None = None) -> Flask:
|
||||
app.extensions["dividend_ledger"] = DividendLedger(
|
||||
app.config.get("DIVIDEND_LEDGER_PATH") or (data_root / "dividends" / "ledger.json")
|
||||
)
|
||||
from .backtest_store import BacktestStore
|
||||
app.extensions["backtest_store"] = BacktestStore(
|
||||
app.config.get("BACKTEST_STORE_PATH") or (data_root / "backtest" / "runs.json")
|
||||
)
|
||||
|
||||
# shared in-process daily cache + app-internal data scheduler (independent of
|
||||
# Hermes — this app runs on its own server).
|
||||
@@ -771,6 +775,88 @@ def create_app(config: dict[str, Any] | None = None) -> Flask:
|
||||
)
|
||||
return jsonify(res.to_dict())
|
||||
|
||||
@app.post("/api/v1/backtest/run")
|
||||
def backtest_run_endpoint():
|
||||
"""Run a strict, event-driven PIT backtest and persist the result.
|
||||
|
||||
Fails closed on readiness: if the requested window (or the derived
|
||||
default window) is not fully PIT-covered, returns 400 listing the
|
||||
missing inputs rather than fabricating an invalid run.
|
||||
"""
|
||||
from .backtest_readiness import evaluate_readiness
|
||||
from .backtest_engine import run_event_backtest
|
||||
from .factor_vintages import FactorVintageStore
|
||||
from .siamchart_vintages import SiamchartVintageStore
|
||||
from .simulation import load_price_snapshot
|
||||
from .pit_scorer import PitScoreProvider, make_pit_score_fn
|
||||
from pathlib import Path as _Path
|
||||
|
||||
body = request.get_json(silent=True) or {}
|
||||
data_root = _Path(__file__).resolve().parents[1] / "data"
|
||||
fstore = FactorVintageStore(data_root)
|
||||
sstore = SiamchartVintageStore(data_root)
|
||||
price_series = load_price_snapshot()
|
||||
|
||||
# readiness gate (starts below fail closed)
|
||||
start = body.get("start")
|
||||
end = body.get("end")
|
||||
ready = evaluate_readiness(
|
||||
factor_store=fstore, siamchart_store=sstore,
|
||||
price_series=price_series, start=start, end=end,
|
||||
)
|
||||
if not ready.ready:
|
||||
return jsonify({
|
||||
"error": "not enough PIT coverage for this backtest window",
|
||||
"missing": ready.missing,
|
||||
"recommended_start": ready.recommended_start,
|
||||
"recommended_end": ready.recommended_end,
|
||||
"coverage": ready.coverage,
|
||||
}), 400
|
||||
start = start or ready.recommended_start
|
||||
end = end or ready.recommended_end
|
||||
assert start is not None and end is not None # readiness guarantees this
|
||||
capital = float(body.get("capital") or 1_000_000)
|
||||
|
||||
# seed siamchart vintage baseline + PIT scorer
|
||||
_seed_siamchart_vintage(sstore)
|
||||
provider = PitScoreProvider(fstore, _load_siamchart_snapshot(),
|
||||
siamchart_store=sstore)
|
||||
score_fn = make_pit_score_fn(provider)
|
||||
|
||||
use_ledger = bool(body.get("use_ledger"))
|
||||
ledger = None
|
||||
if use_ledger:
|
||||
stored = app.extensions.get("dividend_ledger")
|
||||
if stored is not None and stored.symbols():
|
||||
ledger = stored
|
||||
else:
|
||||
from .dividend_ledger import build_dps_ledger
|
||||
ledger = build_dps_ledger(_load_siamchart_snapshot())
|
||||
|
||||
res = run_event_backtest(
|
||||
start=start, end=end, capital=capital,
|
||||
factor_store=fstore, siamchart_store=sstore,
|
||||
dividend_ledger=ledger, price_series=price_series,
|
||||
score_fn=score_fn,
|
||||
)
|
||||
record = res.to_dict()
|
||||
record["use_ledger"] = use_ledger
|
||||
persisted = app.extensions["backtest_store"].add(record)
|
||||
return jsonify(persisted), 201
|
||||
|
||||
@app.get("/api/v1/backtest/run/<run_id>")
|
||||
def backtest_run_get(run_id: str):
|
||||
store = app.extensions["backtest_store"]
|
||||
rec = store.get(run_id)
|
||||
if rec is None:
|
||||
return jsonify({"error": "run not found"}), 404
|
||||
return jsonify(rec)
|
||||
|
||||
@app.get("/api/v1/backtest/run")
|
||||
def backtest_runs_direct():
|
||||
store = app.extensions["backtest_store"]
|
||||
return jsonify({"runs": store.all()})
|
||||
|
||||
@app.post("/api/v1/backtest")
|
||||
def run_backtest_endpoint():
|
||||
"""Run a real backtest over [start, end] with capital; persist result."""
|
||||
|
||||
88
backend/app/backtest_store.py
Normal file
88
backend/app/backtest_store.py
Normal file
@@ -0,0 +1,88 @@
|
||||
"""Durable JSON store for event-driven backtest runs (Task 6).
|
||||
|
||||
The previous backtest run history lived only in process-local memory
|
||||
(``app.extensions.setdefault("backtest_runs", [])``) and was lost on restart.
|
||||
This store persists every event-driven run (with its input coverage, event
|
||||
timeline, and scorer provenance) to an atomic JSON file at ``data/backtest/runs.json``
|
||||
so history survives restarts and can be audited.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import datetime as dt
|
||||
import json
|
||||
import os
|
||||
import threading
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from typing import Any, Optional
|
||||
|
||||
|
||||
def _utcnow_iso() -> str:
|
||||
return dt.datetime.now(dt.timezone.utc).isoformat(timespec="seconds")
|
||||
|
||||
|
||||
class BacktestStoreError(ValueError):
|
||||
pass
|
||||
|
||||
|
||||
class BacktestStore:
|
||||
"""Append-only, thread-safe, atomic JSON store of backtest runs."""
|
||||
|
||||
def __init__(self, path: Optional[Path | str] = None) -> None:
|
||||
self.path = Path(path) if path else None
|
||||
self._lock = threading.Lock()
|
||||
self._runs: list[dict[str, Any]] = []
|
||||
if self.path and self.path.is_file():
|
||||
self._load()
|
||||
|
||||
# -- persistence ------------------------------------------------------
|
||||
def _load(self) -> None:
|
||||
p = self.path
|
||||
if p is None:
|
||||
return
|
||||
try:
|
||||
payload = json.loads(p.read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError) as exc:
|
||||
raise BacktestStoreError(f"cannot load backtest store: {exc}") from exc
|
||||
if isinstance(payload, list):
|
||||
self._runs = payload
|
||||
elif isinstance(payload, dict) and isinstance(payload.get("runs"), list):
|
||||
self._runs = payload["runs"]
|
||||
else:
|
||||
self._runs = []
|
||||
|
||||
def _save(self) -> None:
|
||||
if not self.path:
|
||||
return
|
||||
self.path.parent.mkdir(parents=True, exist_ok=True)
|
||||
tmp = self.path.with_suffix(".tmp")
|
||||
tmp.write_text(
|
||||
json.dumps({"runs": self._runs}, ensure_ascii=False, default=str),
|
||||
encoding="utf-8",
|
||||
)
|
||||
tmp.replace(self.path)
|
||||
|
||||
# -- API --------------------------------------------------------------
|
||||
def add(self, record: dict[str, Any]) -> dict[str, Any]:
|
||||
"""Persist a completed run record; returns it with an id/timestamp."""
|
||||
with self._lock:
|
||||
rec = dict(record)
|
||||
rec.setdefault("id", str(uuid.uuid4())[:12])
|
||||
rec.setdefault("ran_at", _utcnow_iso())
|
||||
self._runs.append(rec)
|
||||
self._save()
|
||||
return rec
|
||||
|
||||
def all(self) -> list[dict[str, Any]]:
|
||||
import copy
|
||||
with self._lock:
|
||||
return copy.deepcopy(self._runs)
|
||||
|
||||
def get(self, run_id: str) -> Optional[dict[str, Any]]:
|
||||
import copy
|
||||
with self._lock:
|
||||
for r in self._runs:
|
||||
if r.get("id") == run_id:
|
||||
return copy.deepcopy(r)
|
||||
return None
|
||||
53
backend/tests/test_backtest_store.py
Normal file
53
backend/tests/test_backtest_store.py
Normal file
@@ -0,0 +1,53 @@
|
||||
"""Tests for the durable backtest run store (Task 6)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
|
||||
from app.backtest_store import BacktestStore
|
||||
|
||||
|
||||
class StoreTest(unittest.TestCase):
|
||||
def tearDown(self):
|
||||
self._dir = getattr(self, "_dir", None)
|
||||
if self._dir:
|
||||
import shutil
|
||||
shutil.rmtree(self._dir, ignore_errors=True)
|
||||
|
||||
def _store(self):
|
||||
self._dir = tempfile.mkdtemp()
|
||||
return BacktestStore(Path(self._dir) / "runs.json")
|
||||
|
||||
def test_add_persists_and_assigns_id(self):
|
||||
s = self._store()
|
||||
rec = s.add({"start": "2026-01-01", "end": "2026-06-01", "final_equity": 123.0})
|
||||
self.assertIn("id", rec)
|
||||
self.assertIn("ran_at", rec)
|
||||
self.assertEqual(len(s.all()), 1)
|
||||
|
||||
def test_survives_restart(self):
|
||||
path = Path(tempfile.mkdtemp()) / "runs.json"
|
||||
s1 = BacktestStore(path)
|
||||
rec = s1.add({"start": "2026-01-01", "end": "2026-06-01"})
|
||||
# new instance reads from disk
|
||||
s2 = BacktestStore(path)
|
||||
self.assertEqual(len(s2.all()), 1)
|
||||
got = s2.get(rec["id"])
|
||||
assert got is not None
|
||||
self.assertEqual(got["start"], "2026-01-01")
|
||||
|
||||
def test_get_returns_none_for_unknown(self):
|
||||
s = self._store()
|
||||
self.assertIsNone(s.get("nope"))
|
||||
|
||||
def test_all_returns_copies(self):
|
||||
s = self._store()
|
||||
s.add({"x": 1})
|
||||
s.all()[0]["x"] = 999
|
||||
self.assertEqual(s.all()[0]["x"], 1)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user