"""
build_flows.py -- [2026-09-28] Per-minute buy / sell volume for every stock-day in the
22-day set, from the tick cache (each trade labeled against the latest quote, <= 2 s old:
at/above ask = buy, at/below bid = sell, else vs the midpoint; no fresh quote -> tick rule).
Saves flow/<date>/<SYM>.json = {"buy": [390], "sell": [390], "other": [390]} (9:30 = index 0).
Skips files already built, so it can be rerun safely.
"""
import gzip
import json
import sys
import time
from datetime import datetime
from pathlib import Path
from zoneinfo import ZoneInfo

HERE = Path(__file__).resolve().parent
CACHE = Path("/var/www/screener/trade/data/simulations/cache")
ET = ZoneInfo("America/New_York")


def flows(path):
    buy, sell, other = [0.0] * 390, [0.0] * 390, [0.0] * 390
    q = None
    last_p, last_side = None, 0
    with open(path) as f:
        for line in f:
            ts, kind, a, b = json.loads(line)
            t = datetime.fromisoformat(ts)
            if kind == "quote":
                q = (t, a, b)
                continue
            p, size = a, b
            side = 0
            if q and (t - q[0]).total_seconds() <= 2 and q[1] > 0 and q[2] >= q[1]:
                bid, ask = q[1], q[2]
                mid = (bid + ask) / 2
                side = 1 if p >= ask else -1 if p <= bid else (1 if p > mid else -1 if p < mid else 0)
            if side == 0 and last_p is not None:
                side = 1 if p > last_p else -1 if p < last_p else last_side
            last_p, last_side = p, side
            te = t.astimezone(ET)
            m = (te.hour - 9) * 60 + te.minute - 30
            if 0 <= m < 390:
                (buy if side > 0 else sell if side < 0 else other)[m] += size
    return {"buy": buy, "sell": sell, "other": other}


if __name__ == "__main__":
    t0 = time.time()
    n = 0
    for p in sorted(HERE.glob("bars/bars_2026-*.json.gz")):
        d = json.load(gzip.open(p, "rt"))
        day = d["date"]
        out = HERE / "flow" / day
        out.mkdir(parents=True, exist_ok=True)
        for s in d["top30_current"]:
            f = out / f"{s}.json"
            if f.exists():
                continue
            src = CACHE / day / f"{s}.jsonl"
            if not src.exists():
                continue
            f.write_text(json.dumps(flows(src)))
            n += 1
        print(f"{day} done ({n} files, {time.time() - t0:.0f}s)", flush=True)
    print("FLOWS_DONE", flush=True)
