"""
history_flows.py -- [2026-10-01] Buy/sell flow per minute for every library stock-day (the user:
add buy/sell flow to the playbook matching). For each day (newest first) and each of its 30
list stocks without flow/<day>/<SYM>.json: download the day's trades + quotes (9:30-16:00) from
Alpaca to a temp file, compute buy / sell volume per minute (build_flows.flows: quote rule,
tick-rule fallback), save the small flow file, delete the ticks. Resumable; stops at 9:10 ET and
never runs 9:10-16:05.
    python3 history_flows.py
    python3 history_flows.py --oldest-first   [2026-10-03] second worker from the other end; stops when it
                                              reaches a day the newest-first worker has started
"""
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
sys.path.insert(0, str(HERE)); sys.path.insert(0, str(HERE.parents[1]))
from alpaca_client import get_client
from simulate import cache_symbol
from build_flows import flows

ET = ZoneInfo("America/New_York")
TMP = Path("/tmp/claude-0/-var-www-screener-trade/5844eef7-c648-46de-a03c-fe0b1434b640/scratchpad/flowtmp")


def market_hours():
    t = datetime.now(ET)
    return t.weekday() < 5 and (9, 10) <= (t.hour, t.minute) < (16, 5)


def main():
    TMP.mkdir(parents=True, exist_ok=True)
    client = get_client()
    t0 = time.time()
    done = 0
    old = "--oldest-first" in sys.argv
    for p in sorted((HERE / "bars").glob("bars_20*.json.gz"), reverse=not old):
        day = p.name[5:15]
        out = HERE / "flow" / day
        if old and out.exists() and not (out / ".oldest").exists():
            print(f"{day}: reached the newest-first worker", flush=True)
            break
        d = json.load(gzip.open(p, "rt"))
        out.mkdir(parents=True, exist_ok=True)
        if old:
            (out / ".oldest").touch()
        todo = [s for s in d["top30_current"] if not (out / f"{s}.json").exists()]
        if not todo:
            continue
        d0 = datetime.fromisoformat(day).replace(tzinfo=ET)
        for s in todo:
            if market_hours():
                print("stopped: market hours", flush=True)
                return
            tmp = TMP / f"{day}_{s}{'_old' if old else ''}.jsonl"
            try:
                cache_symbol(client, s, d0.replace(hour=9, minute=30), d0.replace(hour=16), tmp)
                (out / f"{s}.json").write_text(json.dumps(flows(tmp)))
            except Exception as e:
                print(day, s, "failed", e, flush=True)
            finally:
                tmp.unlink(missing_ok=True)
                Path(str(tmp) + ".tmp").unlink(missing_ok=True) if hasattr(Path, "unlink") else None
        done += 1
        print(f"{day}: {len(todo)} stocks [{(time.time() - t0) / 60:.0f} min, {done} days]", flush=True)
    print("HISTORY_FLOWS_DONE", flush=True)


if __name__ == "__main__":
    main()
