#!/usr/bin/env python3
"""
NPCI product statistics backfill fetcher.

Fetches every JSON table behind npci.org.in product statistics pages via the
site's Strapi API (in-page fetch from a real Chromium session, so Akamai
bot-protection passes). Layout:

  raw/npci/<product>/<stream>/<year_range|mYYYY-MM>.json

API contract (reverse-engineered 2026-09-18):
  - tab catalog:        GET /api/product-statistic/tabs?locale=en
  - FY monthly tables:  tab/detail?product_name=<p>&tab_name=<slug>&excel_type=monthly&year_range=YYYY-YY
  - month-scoped tables:tab/detail?...&excel_type=daily|uptime|downtime|monthly&month=<Month>&year=YYYY
    (for mandates excel_type=monthly + tab_name=mandate-creation|execution)
  - Responses: {"status":200,"data":{"results":[...],"totalCount":N}}

Usage:
  python3 npci_fetch.py                 # full catalog
  python3 npci_fetch.py --product upi   # one product
  python3 npci_fetch.py --catalog-only
"""

import argparse
import asyncio
import calendar
import datetime
import json
import sys
from pathlib import Path

from playwright.async_api import async_playwright

BASE = "https://www.npci.org.in"
API = f"{BASE}/api/product-statistic"
RAW = Path("/home/workspace/Projects/payments-stat-hub/raw/npci")
MANIFEST = RAW / "manifest.jsonl"

PAGE_TIMEOUT = 30000
SLEEP_BETWEEN = 0.35
PAGE_SIZE = 100
STOP_EMPTY_MONTHS = 12          # stop month-loop after this many consecutive empty months
FLOOR = (2016, 1)               # never probe before Jan 2016


def log(msg):
    print(f"{datetime.datetime.now().strftime('%H:%M:%S')} {msg}", flush=True)


class NpciSession:
    def __init__(self, browser):
        self.browser = browser
        self.page = None

    async def start(self):
        self.page = await self.browser.new_page()
        resp = await self.page.goto(f"{BASE}/product/upi/product-statistics",
                                    wait_until="domcontentloaded", timeout=PAGE_TIMEOUT)
        await self.page.wait_for_timeout(5000)
        if resp is not None and resp.status >= 400:
            raise RuntimeError(f"NPCI page load blocked: HTTP {resp.status}")
        for attempt in range(3):
            status, _ = await self.api("tabs", {"locale": "en"})
            if status == 200:
                return
            log(f"API probe got {status}, refreshing session (attempt {attempt + 1})")
            await self.refresh()
        raise RuntimeError("NPCI API not reachable after session refreshes")

    async def refresh(self):
        await self.page.goto(f"{BASE}/product/upi/product-statistics",
                             wait_until="domcontentloaded", timeout=PAGE_TIMEOUT)
        await self.page.wait_for_timeout(4000)

    async def api(self, path, params):
        query = "&".join(f"{k}={v}" for k, v in params.items() if v is not None)
        url = f"{API}/{path}?{query}"
        expr = ("(async () => { try { const r = await fetch(" + json.dumps(url) +
                ", {headers: {'Accept': 'application/json'}});"
                " return [r.status, await r.text()]; }"
                " catch (e) { return [-1, String(e)]; } })()")
        try:
            result = await self.page.evaluate(expr)
        except Exception as e:
            return -1, f"evaluate-failed: {e}"
        status, text = result
        return status, text


async def fetch_detail(sess, product, tab_name, excel_type, year_range=None,
                       month=None, year=None):
    """Fetch one table fully (all pages). Returns (rows, total, status)."""
    all_rows, status = [], -1
    page_no = 1
    data = {}
    while True:
        params = {"product_name": product.lower(), "tab_name": tab_name,
                  "excel_type": excel_type, "page_no": page_no,
                  "page_size": PAGE_SIZE, "locale": "en"}
        if year_range:
            params["year_range"] = year_range
        if month:
            params["month"] = month
        if year:
            params["year"] = year
        status, text = await sess.api("tab/detail", params)
        if status != 200:
            break
        try:
            body = json.loads(text)
        except json.JSONDecodeError:
            break
        data = body.get("data") or {}
        results = data.get("results") or []
        total = data.get("totalCount") or 0
        all_rows.extend(results)
        if not results or len(all_rows) >= total:
            break
        page_no += 1
        await asyncio.sleep(SLEEP_BETWEEN)
    return all_rows, (data.get("totalCount") if status == 200 else None), status


def save_raw(product, stream, key, status, text):
    d = RAW / product.lower() / stream
    d.mkdir(parents=True, exist_ok=True)
    path = d / f"{key}.json"
    path.write_text(text)
    return path


def manifest_line(url, path, status, total, rows):
    rec = {"ts": datetime.datetime.now().isoformat(timespec="seconds"),
           "url": url, "file": str(path), "status": status,
           "totalCount": total, "rows": rows}
    with open(MANIFEST, "a") as f:
        f.write(json.dumps(rec, ensure_ascii=False) + "\n")


async def fetch_catalog(sess):
    status, text = await sess.api("tabs", {"locale": "en"})
    if status != 200:
        raise RuntimeError(f"tabs catalog fetch failed: {status}")
    (RAW / "tabs-catalog.json").write_text(text)
    body = json.loads(text)
    out = []
    for p in body["data"][0]["products"]:
        for t in p.get("tab_details") or []:
            out.append({
                "product": p["product_name"],
                "tab_name": t["tab_name"],
                "slug": t["slug"],
                "tab_type": t.get("tab_type"),
                "years": [y["year_range"] for y in (t.get("year_range") or [])],
            })
    return out


def classify_streams(tabs):
    """Group month-scoped tabs of one product into deduped fetch streams."""
    streams = {}
    for t in tabs:
        slug = t["slug"]
        if t["tab_type"] == "pdf":
            continue
        if t["years"]:
            continue
        if slug == "mandate-creation":
            streams["mandate-creation"] = {"tab_name": slug, "excel": ["monthly"], "kind": "mandate"}
        elif slug == "mandate-execution":
            streams["mandate-execution"] = {"tab_name": slug, "excel": ["monthly"], "kind": "mandate"}
        elif "daily" in slug:
            streams.setdefault("daily", {"tab_name": slug, "excel": ["daily"], "kind": "generic"})
        elif "uptime" in slug:
            streams.setdefault("uptime", {"tab_name": slug, "excel": ["uptime"], "kind": "generic"})
        elif "downtime" in slug:
            streams.setdefault("downtime", {"tab_name": slug, "excel": ["downtime"], "kind": "generic"})
        else:
            streams[slug] = {"tab_name": slug,
                             "excel": ["daily", "uptime", "downtime", "monthly"],
                             "kind": "probe"}
    return streams


def build_plan(catalog):
    plan = {"fy": {}, "month": {}}
    for t in catalog:
        prod = t["product"]
        if t["tab_type"] == "pdf":
            continue
        if t["years"]:
            for fy in t["years"]:
                plan["fy"].setdefault(prod, []).append({"key": fy, "tab_name": t["slug"]})
        else:
            pass  # handled below per product
    for t in catalog:
        prod = t["product"]
        if t["tab_type"] == "pdf" or t["years"]:
            continue
        if prod not in plan["month"]:
            plan["month"][prod] = classify_streams([x for x in catalog
                                                    if x["product"] == prod
                                                    and x["tab_type"] != "pdf"
                                                    and not x["years"]])
    return plan


def months_desc():
    today = datetime.date.today()
    y, m = today.year, today.month
    out = []
    while (y, m) >= FLOOR:
        out.append((y, calendar.month_name[m]))
        m -= 1
        if m == 0:
            m = 12
            y -= 1
    return out


async def fetch_fy_units(sess, prod, units):
    ok = empty = fail = 0
    for u in units:
        path = RAW / prod.lower() / "product-monthly" / f"{u['key']}.json"
        if path.exists() and "results" in path.read_text():
            ok += 1
            continue
        rows, total, status = await fetch_detail(
            sess, prod, u["tab_name"], "monthly", year_range=u["key"])
        text = json.dumps({"data": {"results": rows, "totalCount": total}},
                          ensure_ascii=False) if rows else \
               json.dumps({"empty": True, "last_status": status})
        p = save_raw(prod, "product-monthly", u["key"], status, text)
        manifest_line(f"fy:{prod}:{u['tab_name']}:{u['key']}", p, status, total, len(rows))
        if rows:
            ok += 1
        elif status == 200:
            empty += 1
        else:
            fail += 1
        log(f"  fy {prod} {u['key']} status={status} rows={len(rows)}")
        await asyncio.sleep(SLEEP_BETWEEN)
    return ok, empty, fail


def has_rows(path):
    try:
        body = json.loads(path.read_text())
        return bool(body.get("data", {}).get("results"))
    except Exception:
        return False


def row_sig(rows):
    if not rows:
        return None
    return "|".join(sorted(rows[0].keys()))


async def fetch_month_stream(sess, prod, stream_key, stream, counts, seen_sigs):
    """Backward month loop for one stream, with dedupe + idempotent resume."""
    ok = empty = fail = 0
    consecutive_empty = 0
    got_success = False
    pinned_excel = stream["excel"][0]
    tried_pin = False
    d = RAW / prod.lower() / stream_key
    d.mkdir(parents=True, exist_ok=True)

    for (yy, mm) in months_desc():
        key = f"{yy}-{mm}"
        path = d / f"{key}.json"
        if path.exists() and has_rows(path):
            ok += 1
            consecutive_empty = 0
            continue
        KNOWN_KINDS = ("daily", "uptime", "downtime", "mandate")
        if not tried_pin and stream["kind"] not in KNOWN_KINDS:
            # pin excel_type on first (candidate x recent month) that returns rows
            pinned = None
            sig = None
            recent = months_desc()[:4]
            for (py, pm) in recent:
                for cand in stream["excel"]:
                    rows, total, status = await fetch_detail(
                        sess, prod, stream["tab_name"], cand, month=pm, year=py)
                    await asyncio.sleep(SLEEP_BETWEEN)
                    if status == 200 and rows:
                        pinned = cand
                        sig = row_sig(rows)
                        break
                if pinned:
                    if sig and sig in seen_sigs.get(prod, set()) and stream["kind"] == "probe":
                        log(f"  {prod}/{stream_key}: duplicate of earlier stream (keys {sig}), skipping")
                        return ok, empty, fail + 1
                    break
            if pinned is None:
                path.write_text(json.dumps({"empty": True, "probe": True}))
                manifest_line(f"probe:{prod}:{stream_key}", path, status, None, 0)
                empty += 1
                log(f"  {prod}/{stream_key}: no excel_type matched, giving up")
                return ok, empty, fail + 1
            pinned_excel = pinned
            tried_pin = True
            seen_sigs.setdefault(prod, set()).add(sig)
        elif not tried_pin:
            pinned_excel = stream["excel"][0]
            tried_pin = True
        rows, total, status = await fetch_detail(
            sess, prod, stream["tab_name"], pinned_excel, month=mm, year=yy)
        text = json.dumps({"data": {"results": rows, "totalCount": total},
                           "excel_type": pinned_excel}, ensure_ascii=False) if rows else \
               json.dumps({"empty": True, "last_status": status})
        p = save_raw(prod, stream_key, key, status, text)
        manifest_line(f"m:{prod}:{stream_key}:{pinned_excel}:{key}", p, status, total, len(rows))
        if rows:
            ok += 1
            got_success = True
            consecutive_empty = 0
        elif status == 200:
            empty += 1
            consecutive_empty += 1
        else:
            fail += 1
            consecutive_empty += 1
        if consecutive_empty >= STOP_EMPTY_MONTHS and got_success:
            log(f"  {prod}/{stream_key}: {STOP_EMPTY_MONTHS} empty months, stopping loop")
            break
        await asyncio.sleep(SLEEP_BETWEEN)
    return ok, empty, fail


async def main():
    ap = argparse.ArgumentParser()
    ap.add_argument("--product", help="fetch one product only (e.g. upi)")
    ap.add_argument("--catalog-only", action="store_true")
    args = ap.parse_args()

    async with async_playwright() as pw:
        browser = await pw.chromium.launch(headless=True)
        sess = NpciSession(browser)
        await sess.start()
        log("session up, API verified")

        catalog = await fetch_catalog(sess)
        log(f"catalog: {len(catalog)} tabs")
        if args.catalog_only:
            await browser.close()
            return

        plan = build_plan(catalog)
        if args.product:
            prod = args.product.lower()
            plan["fy"] = {k: v for k, v in plan["fy"].items() if k.lower() == prod}
            plan["month"] = {k: v for k, v in plan["month"].items() if k.lower() == prod}

        totals = {"ok": 0, "empty": 0, "fail": 0}
        seen_sigs = {}
        for prod, units in plan["fy"].items():
            log(f"[fy] {prod}: {len(units)} FY units")
            ok, empty, fail = await fetch_fy_units(sess, prod, units)
            totals["ok"] += ok
            totals["empty"] += empty
            totals["fail"] += fail

        for prod, streams in plan["month"].items():
            for stream_key, stream in streams.items():
                log(f"[month] {prod}/{stream_key} ({stream['kind']})")
                ok, empty, fail = await fetch_month_stream(
                    sess, prod, stream_key, stream, totals, seen_sigs)
                totals["ok"] += ok
                totals["empty"] += empty
                totals["fail"] += fail

        await browser.close()
        log(f"DONE ok={totals['ok']} empty={totals['empty']} fail={totals['fail']}")


if __name__ == "__main__":
    asyncio.run(main())
