File size: 5,769 Bytes
e317359
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
8c1bce7
e317359
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
#!/usr/bin/env python3
"""Real, isolated 5-minute Binance forecast/observation/scoring acceptance.

Issue near the beginning of a five-minute interval; resolve after the following
interval closes. This uses actual model APIs and wall time, never a mocked clock.
Its horizon differs from the leaderboard and its scores must remain isolated.
"""
from __future__ import annotations

import argparse
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timezone
import hashlib
import json
from pathlib import Path
import sys
import time
from types import SimpleNamespace

ROOT = Path(__file__).resolve().parents[1]
sys.path[:0] = [str(ROOT), str(ROOT / "src")]

import numpy as np
import pandas as pd
import requests
from dotenv import load_dotenv
from scripts.run_online_eval import load_model_specs, make_predictor, model_display_name, model_output_slug
from tsfm_bench.eval.prequential import PrequentialStore, forecast_deadline


class Source:
    _settings = SimpleNamespace(prediction_length=1)

    def __init__(self, rows):
        self.rows = rows

    def get_metadata(self, name):
        return SimpleNamespace(domain="Finance", frequency="5min")

    def stream(self, name):
        yield SimpleNamespace(target=np.asarray([float(row[4]) for row in self.rows]),
                              start=str(pd.Timestamp(self.rows[0][0], unit="ms")), freq="5min")


def main():
    load_dotenv(ROOT / ".env")
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("phase", choices=["issue", "resolve"])
    parser.add_argument("--output-root", type=Path, default=ROOT / "outputs/live-acceptance")
    args = parser.parse_args()
    if args.output_root.resolve().is_relative_to((ROOT / "space/results").resolve()):
        raise SystemExit("Acceptance results must stay outside production results")
    args.output_root.mkdir(parents=True, exist_ok=True)
    if args.phase == "issue":
        remaining = 300 - time.time() % 300
        if remaining < 180:
            print(f"Waiting {remaining + 3:.0f}s for a full forecast issue window", flush=True)
            time.sleep(remaining + 3)
    response = requests.get("https://data-api.binance.vision/api/v3/klines", timeout=45,
                            params={"symbol": "BTCUSDT", "interval": "5m", "limit": 130})
    response.raise_for_status()
    rows = response.json()
    fetched = datetime.now(timezone.utc).isoformat()
    (args.output_root / f"{args.phase}-observations.json").write_text(json.dumps({
        "fetched_at": fetched, "provider": response.url, "rows": rows}, indent=2))
    store = PrequentialStore(args.output_root)
    specs = [s for s in load_model_specs(ROOT / "configs/models/online_tsfm.yaml") if s.get("enabled", True)]
    if args.phase == "issue":
        if store.tasks():
            raise SystemExit("Use a new output directory for a new acceptance window")
        periods = pd.period_range(pd.Timestamp(rows[0][0], unit="ms"), periods=len(rows), freq="5min")
        task, _ = store.ensure_task({
            "source_task_name": "binance_acceptance/task", "dataset": "binance_acceptance/5min/short",
            "domain": "Finance", "frequency": "5min", "prediction_length": 1,
            "context": [float(row[4]) for row in rows],
            "context_timestamps": [p.start_time.isoformat() for p in periods],
            "context_start": str(periods[0]), "context_end": str(periods[-1]), "data_fetched_at": fetched,
        })
        print("Forecast deadline:", forecast_deadline(task), flush=True)

        def issue(spec):
            name = model_display_name(spec)
            try:
                ok = store.issue_forecast(task, spec, model_name=name, model_slug=model_output_slug(spec),
                                          predictor=make_predictor(spec, 1, False))
                return {"model": name, "frozen": ok}
            except Exception as exc:
                return {"model": name, "frozen": False, "error": str(exc)}

        results = []
        with ThreadPoolExecutor(max_workers=8) as pool:
            for future in as_completed([pool.submit(issue, spec) for spec in specs]):
                result = future.result()
                results.append(result)
                print(json.dumps(result), flush=True)
        hashes = {p.name: hashlib.sha256(p.read_bytes()).hexdigest()
                  for p in (store.forecasts_dir / task["task_id"]).glob("*.json")}
        report = {"phase": "issue", "isolated": True, "task_id": task["task_id"],
                  "forecast_deadline": forecast_deadline(task).isoformat(), "models": results,
                  "frozen_file_sha256": hashes}
        (args.output_root / "acceptance-issue.json").write_text(json.dumps(report, indent=2))
        return int(not all(r["frozen"] for r in results))
    report = json.loads((args.output_root / "acceptance-issue.json").read_text())
    for name, digest in report["frozen_file_sha256"].items():
        assert hashlib.sha256((store.forecasts_dir / report["task_id"] / name).read_bytes()).hexdigest() == digest
    cycle = store.resolve_ready(Source(rows), ["binance_acceptance/task"])
    task = store.tasks()[0]
    report = {"phase": "resolve", "isolated": True, "frozen_hashes_unchanged": True,
              "observed_at": fetched, "task_status": task["status"],
              "scored_models": task["scored_models"], "expected_models": len(specs),
              "pending_tasks": cycle.pending_tasks}
    (args.output_root / "acceptance-resolve.json").write_text(json.dumps(report, indent=2))
    print(json.dumps(report, indent=2))
    return int(task["status"] != "resolved" or len(task["scored_models"]) != len(specs))


if __name__ == "__main__":
    raise SystemExit(main())