#!/usr/bin/env python3
"""Build all hub data exports + catalog from payments.duckdb.

Run: python3 Skills/payments-fetch/scripts/build_hub_data.py
Outputs to Projects/payments-stat-hub/data/ (summary/monthly/daily/quality/psi/
bankwise/atmpos/catalog json+csv). Every catalog entry carries its source URL,
fetch script and coverage so the hub can show provenance for each dataset.
"""
import json
import re
from datetime import date
from pathlib import Path

import duckdb

ROOT = Path("/home/workspace/Projects/payments-stat-hub")
OUT = ROOT / "data"
OUT.mkdir(parents=True, exist_ok=True)
con = duckdb.connect(str(ROOT / "data" / "payments.duckdb"))

NPCI_ECO_URL = "https://www.npci.org.in/product/upi/ecosystem-statistics"
NPCI_STATS = "https://www.npci.org.in/product/{p}/product-statistics"
RBI_PSI = "https://www.rbi.org.in/Scripts/PSIUserView.aspx"
RBI_BANKWISE = "https://www.rbi.org.in/Scripts/BankwiseCurrency.aspx"
RBI_ATM = "https://www.rbi.org.in/Scripts/ATMView.aspx"
ZOPUB = "https://zo.pub/cashlessconsumer/payments-stat-hub"

PRODUCT_TITLES = {
    "upi": "UPI (Unified Payments Interface)",
    "imps": "IMPS (Immediate Payment Service)",
    "netc": "NETC FASTag",
    "bhim": "BHIM (UPI app)",
    "nfs": "NFS (National Financial Switch ATM)",
    "aeps": "AePS (Aadhaar enabled Payment System)",
    "cts": "CTS (Cheque Truncation System)",
    "nach": "NACH (National Automated Clearing House)",
    "e-rupi": "e-RUPI vouchers",
}

def q(sql):
    return con.execute(sql).fetchall()

def cols(table):
    return [r[0] for r in con.execute(f"DESCRIBE {table}").fetchall()]

def month_minmax(table, col="month"):
    r = con.execute(f"SELECT min({col}), max({col}), count(*) FROM {table}").fetchone()
    return {"from": str(r[0]), "to": str(r[1]), "rows": r[2]}

def dump(name, obj):
    p = OUT / name
    p.write_text(json.dumps(obj, ensure_ascii=False, indent=1, default=str))
    print(f"  {name} ({p.stat().st_size//1024} KB)")

# ---- 1. unified monthly -----------------------------------------------------
con.execute("""CREATE OR REPLACE VIEW all_monthly AS
SELECT product, month, volume_mn, value_cr, banks FROM npci_product_monthly
UNION ALL
SELECT product, month, total_volume_mn, total_value_cr, NULL FROM npci_trended_monthly""")
monthly = q("SELECT * FROM all_monthly ORDER BY product, month")
dump("monthly.json", [{"product": p, "month": str(m), "volume_mn": v, "value_cr": c, "banks": b}
                      for p, m, v, c, b in monthly])
con.execute("COPY (SELECT * FROM all_monthly ORDER BY month, product) TO '"
            + str(OUT / "monthly.csv") + "' (HEADER)")

# ---- 2. summary -------------------------------------------------------------
mtd = q("""SELECT product, volume_mn, value_cr FROM npci_daily
WHERE day = (SELECT max(day) FROM npci_daily)""")
mtd_map = {}
for p, v, c in mtd:
    d = mtd_map.setdefault(p, {"volume_mn": 0.0, "value_cr": 0.0})
    d["volume_mn"] += v or 0
    d["value_cr"] += c or 0
prev = {}
cur = {}
for p, m, v, c, b in monthly:
    if v is None:
        continue
    prev[p] = cur.get(p)
    cur[p] = {"month": str(m), "volume_mn": v, "value_cr": c, "banks": b}
summary = []
for p, c in cur.items():
    pv = prev.get(p)
    summary.append({
        "product": p, "title": PRODUCT_TITLES.get(p, p.upper()), "month": c["month"],
        "volume_mn": c["volume_mn"], "value_cr": c["value_cr"], "banks": c["banks"],
        "mom_volume_pct": round((c["volume_mn"] / pv["volume_mn"] - 1) * 100, 1) if pv else None,
        "mtd": mtd_map.get(p) if c["month"][:7] >= sorted(mtd_map.keys() or ["0"])[0][:7] else None,
    })
dump("summary.json", {"generated": str(date.today()), "latest": summary})

# ---- 3. quality (uptime + bank incident/downtime) ---------------------------
quality = {}
for prod, _, period, rj in q("""
    SELECT product, regexp_extract(src_file, '^[^/]+/([^/]+)/', 1), period, rows
    FROM npci_other
    WHERE regexp_extract(src_file, '^[^/]+/([^/]+)/', 1) IN ('uptime','downtime')"""):
    entries = quality.setdefault(prod, {})
    e = entries.setdefault(str(period), {"uptime_pct": None, "unscheduled_downtime": None,
                                         "sys_incidents": None, "bank_incidents": 0,
                                         "downtime_hours": 0.0})
    for r in json.loads(rj) if isinstance(rj, str) else rj:
        if "npci_uptime_for_upi" in r:
            m = re.search(r"[\d.]+", str(r.get("npci_uptime_for_upi", "")))
            e["uptime_pct"] = float(m.group()) if m else None
            inc = str(r.get("no_of_incidents", "")).strip().upper()
            e["sys_incidents"] = None if inc in ("NIL", "") else int(re.sub(r"\D", "", inc) or 0)
            ud = str(r.get("unscheduled_downtime", "")).strip().upper()
            e["unscheduled_downtime"] = None if ud in ("NIL", "") else ud
        else:
            e["bank_incidents"] += int(re.sub(r"[^\d.]", "", str(r.get("incident_count"))) or 0) if re.search(r"\d", str(r.get("incident_count", ""))) else 0
            dts = str(r.get("downtime_in_hours", "")).strip().upper()
            if ":" in dts and re.search(r"\d", dts):
                hh, mm = (dts.split(":") + ["0"])[:2]
                e["downtime_hours"] += (int(re.sub(r"\D", "", hh) or 0)
                                        + int(re.sub(r"\D", "", mm) or 0) / 60)
            elif re.search(r"\d", dts):
                e["downtime_hours"] += float(re.sub(r"[^\d.]", "", dts) or 0)
quality_out = {p: sorted(m.items()) for p, m in quality.items()}
dump("quality.json", quality_out)
con.execute("DROP TABLE IF EXISTS hub_quality")

# ---- 4. psi indicators ------------------------------------------------------
psi_ind = q("""SELECT part, section, indicator, metric, series_type,
                      min(month), max(month), count(*)
               FROM rbi_psi
               WHERE series_type IS NOT NULL
               GROUP BY 1,2,3,4,5 ORDER BY 1,2,3,4,5""")
dump("psi_indicators.json", [{"part": p, "section": sec, "indicator": i,
                              "metric": mt, "series": st, "from": str(a),
                              "to": str(b), "rows": n}
                             for p, sec, i, mt, st, a, b, n in psi_ind])
con.execute("COPY (SELECT * FROM rbi_psi ORDER BY month) TO '" + str(OUT / "rbi_psi.csv") + "' (HEADER)")

psi_series = q("""SELECT indicator, metric || ' | ' || series_type AS key,
                         month, value FROM rbi_psi
                  WHERE series_type IS NOT NULL ORDER BY 1,2,3""")
series = {}
for i, ch, m, v in psi_series:
    series.setdefault(f"{i}||{ch}", []).append([str(m)[:7], v])
dump("psi_series.json", series)

# ---- 5. bankwise totals -----------------------------------------------------
bt = q("""SELECT rail, month, sum(inward_txn), sum(outward_txn)
          FROM rbi_bankwise WHERE bank IS NOT NULL GROUP BY 1,2 ORDER BY 1,2""")
dump("bankwise_totals.json", [{"rail": r, "month": str(m), "in_txn": i, "out_txn": o}
                              for r, m, i, o in bt])
con.execute("COPY (SELECT * FROM rbi_bankwise WHERE bank IS NOT NULL) TO '"
            + str(OUT / "rbi_bankwise.csv") + "' (HEADER)")

# ---- 6. atmpos metrics ------------------------------------------------------
ap = con.execute("""SELECT DISTINCT json_keys(infra) FROM rbi_atmposcard
                    WHERE infra IS NOT NULL LIMIT 1""").fetchall()
metrics = [k for ks in ap for k in ks[0]]
months_ap = con.execute("SELECT min(month), max(month), count(*) FROM rbi_atmposcard").fetchone()
dump("atmpos_metrics.json", {"metrics": sorted(metrics),
                             "from": str(months_ap[0]), "to": str(months_ap[1]),
                             "rows": months_ap[2]})
con.execute("COPY (SELECT * FROM rbi_atmposcard) TO '" + str(OUT / "rbi_atmposcard.csv") + "' (HEADER)")

# ---- 6b. NPCI ecosystem (UPI deep-dive) -------------------------------------
ECO = [
    ("apps", "eco_upi_apps", ["month", "app", "cit_vol_mn", "cit_val_cr", "b2c_vol_mn", "b2c_val_cr",
                               "b2b_vol_mn", "b2b_val_cr", "onus_vol_mn", "onus_val_cr",
                               "total_vol_mn", "total_val_cr"]),
    ("apps-history", "eco_upi_apps_all", None),
    ("bd-td-banks", "eco_top50_member", ["month", "side", "bank", "total_vol_mn", "approved_pct",
                                          "bd_pct", "td_pct", "debit_rev_mn"]),
    ("bd-td-psps", "eco_psp", ["month", "side", "psp", "total_vol_mn", "approved_pct", "bd_pct", "td_pct"]),
    ("statewise", "eco_statewise", ["month", "state", "vol_mn", "vol_pct", "val_cr", "val_pct"]),
    ("mcc", "eco_mcc", ["month", "type", "category", "vol_mn", "vol_pct", "val_cr", "val_pct"]),
    ("p2p-p2m", "eco_p2p_p2m", ["month", "p2p_vol_mn", "p2p_val_cr", "p2m_vol_mn", "p2m_val_cr"]),
    ("top-banks-volval", "eco_top50_volval", ["month", "bank", "vol_mn", "val_cr"]),
    ("chargeback", "eco_chargeback", ["month", "code", "cb_ratio", "cb_received", "cb_represented", "cb_accepted"]),
]
eco_out = {}
eco_meta = {}
for slug, table, fields in ECO:
    try:
        if fields:
            rows = con.execute(f"SELECT {','.join(fields)} FROM {table} ORDER BY month").fetchall()
            recs = [dict(zip(fields, r)) for r in rows]
        else:
            rows = con.execute(f"SELECT * FROM {table} ORDER BY month, app").fetchall()
            names = [d[0] for d in con.execute(f"SELECT * FROM {table} LIMIT 0").description]
            recs = [dict(zip(names, r)) for r in rows]
    except Exception:
        recs = []
    if recs:
        recs = [{k: (None if isinstance(v, float) and v != v else v) for k, v in r.items()} for r in recs]
        eco_out[slug] = recs
        mmx = con.execute(f"SELECT min(month), max(month), count(*) FROM {table}").fetchone()
        eco_meta[slug] = {"from": str(mmx[0]), "to": str(mmx[1]), "rows": mmx[2]}
dump("eco.json", eco_out)

# ---- 7. catalog -------------------------------------------------------------
cat = []
def add(id_, title, desc, kind, src_url, script, coverage, extra=None):
    e = {"id": id_, "title": title, "description": desc, "kind": kind,
         "source": src_url, "fetch_script": script, "coverage": coverage,
         "published_archive": ZOPUB}
    if extra:
        e.update(extra)
    cat.append(e)

for p in ["upi", "imps", "netc", "bhim"]:
    mm = con.execute("""SELECT min(month), max(month), count(*) FROM all_monthly
                        WHERE product=?""", [p]).fetchone()
    if mm[2]:
        add(f"npci-monthly-{p}", f"{PRODUCT_TITLES[p]} — monthly volumes",
            f"Monthly transaction volume (mn) and value (₹ crore) for {PRODUCT_TITLES[p]}, "
            "as published in NPCI product statistics tables.",
            "table+chart", NPCI_STATS.format(p=p), "npci_fetch.py (product-monthly)",
            {"from": str(mm[0]), "to": str(mm[1]), "rows": mm[2]})

for p in ["upi", "nfs"]:
    mm = con.execute("SELECT min(day), max(day), count(*) FROM npci_daily WHERE product=?",
                     [p]).fetchone()
    if mm[2]:
        add(f"npci-daily-{p}", f"{PRODUCT_TITLES[p]} — daily volumes",
            f"Daily transaction volumes for {PRODUCT_TITLES[p]}. Note: NPCI daily tables "
            "restate whole months; later days overwrite earlier values.",
            "table+chart", NPCI_STATS.format(p=p), "npci_fetch.py (daily)",
            {"from": str(mm[0]), "to": str(mm[1]), "rows": mm[2]})

for p in sorted(quality):
    mm = con.execute("""SELECT min(period), max(period), count(*) FROM npci_other
        WHERE product=? AND regexp_extract(src_file, '^[^/]+/([^/]+)/', 1) IN ('uptime','downtime')""",
        [p]).fetchone()
    add(f"quality-{p}", f"{PRODUCT_TITLES[p]} — uptime & declines (incidents)",
        f"Monthly system uptime % and bank-wise outage incidents (count + downtime hours) "
        f"for {PRODUCT_TITLES[p]}. 'Business declines' during outages surface here — "
        "NPCI publishes incidents and downtime per member bank per month.",
        "table+chart", NPCI_STATS.format(p=p), "npci_fetch.py (uptime/downtime streams)",
        {"from": str(mm[0]), "to": str(mm[1]), "rows": sum(len(v) for v in quality[p].values())})

add("rbi-psi", "RBI Payment System Indicators",
    "Monthly indicators from RBI: card additions, ATM/PoS deployments, BHIM Aadhaar Pay, "
    "authorised bank-wise payment volumes and more (long format: indicator × month × value).",
    "table", RBI_PSI, "rbi_fetch.py (psi postbacks)",
    {"from": month_minmax("rbi_psi")["from"], "to": month_minmax("rbi_psi")["to"],
     "rows": month_minmax("rbi_psi")["rows"], "indicators": len(psi_ind)})

add("rbi-bankwise", "RBI bank-wise payment volumes (ECS/NEFT/RTGS/Mobile)",
    "Bank × month transaction counts and values for ECS/NEFT/RTGS and Mobile banking "
    "(inward + outward). RBI's bank-wise UPI pages were discontinued after Jan 2019.",
    "table", RBI_BANKWISE, "rbi_fetch.py (bankwise-volumes postbacks)",
    {"from": month_minmax("rbi_bankwise")["from"], "to": month_minmax("rbi_bankwise")["to"],
     "rows": month_minmax("rbi_bankwise")["rows"]})

add("rbi-atmposcard", "RBI bank-wise ATM / PoS / card statistics",
    "Bank × month infrastructure: ATMs, PoS terminals, card-on-PoS counts, micro-ATMs. "
    "Pre-2017 files are legacy .xls (BIFF).",
    "table", RBI_ATM, "rbi_fetch.py (atm-pos-card postbacks)",
    {"from": month_minmax("rbi_atmposcard")["from"], "to": month_minmax("rbi_atmposcard")["to"],
     "rows": month_minmax("rbi_atmposcard")["rows"]})

add("rbi-upi-bankwise-historic", "RBI bank-wise UPI (historic, to Jan 2019)",
    "Bank-wise UPI volumes from RBI's discontinued table — inside rbi_bankwise dataset, "
    "rail='UPI'. Historic value only; use NPCI for current data.",
    "table", RBI_BANKWISE, "rbi_fetch.py (bankwise-volumes postbacks)",
    {"from": "2017-01", "to": "2019-01", "rows": None, "note": "subset of rbi-bankwise"})

add("npci-warehouse", "Full DuckDB warehouse",
    "Complete warehouse: every table above plus member-bank lists (UPI/NFS direct, sub-member, "
    "RBB, WATM), e-RUPI vouchers, CTS/NACH/AEPS monthly, fetch log with source-file hash per row.",
    "download", ZOPUB, "build_hub_data.py",
    {"file": "payments.duckdb", "rows": sum(
        con.execute(f"SELECT count(*) FROM {t}").fetchone()[0]
        for t in ["npci_product_monthly", "npci_trended_monthly", "npci_daily",
                  "npci_other", "rbi_psi", "rbi_bankwise", "rbi_atmposcard", "fetch_log"])})

# aliases power the search box ("business declines on upi across months")
ALIASES = {
    "quality-upi": "declines decline failed failure outage downtime downtime hours incidents "
                   "uptime availability business technical bank wise member across months "
                   "disruption unstabilised unscheduled",
    "quality-imps": "declines decline outage downtime incidents imps availability months",
    "quality-nfs": "declines decline outage downtime incidents nfs atm availability months",
    "quality-aeps": "declines decline outage downtime incidents aeps aadhaar availability months",
    "quality-netc": "declines decline outage downtime incidents fastag netc toll availability",
    "npci-monthly-upi": "upi monthly volume value banks tps trend growth lakh crore",
    "npci-monthly-imps": "imps monthly volume value trend",
    "npci-monthly-netc": "netc fastag monthly volume toll trend",
    "npci-monthly-bhim": "bhim app monthly volume trend",
    "npci-daily-upi": "upi daily volume trend mtd month to date",
    "npci-daily-nfs": "nfs daily volume atm trend",
    "rbi-psi": "rbi indicators cards pos atm issuance authorised banks aadhaar bhim pay",
    "rbi-bankwise": "bank wise neft rtgs ecs mobile banking volumes per bank",
    "rbi-atmposcard": "atm pos point of sale terminals cards micro atm deployment per bank",
    "rbi-upi-bankwise-historic": "bank wise upi historic discontinued 2019 per bank",
    "npci-warehouse": "download duckdb warehouse all data everything csv export",
    "eco-upi-apps": "upi app tpap cit b2c b2b onus volume value per app",
    "eco-chargeback": "chargeback per beneficiary bank code ratio representment",
    "eco-top50-member": "top remitter beneficiary banks approved bd td decline percentages",
    "eco-psp": "upi payer payee psp performance approved bd td declines",
    "eco-mcc": "merchant category code high transacting categories volume value contribution",
    "eco-p2p-p2m": "upi p2p p2m person to person merchant split volume value",
    "eco-statewise": "upi state wise volume value contribution states",
    "eco-top50-volval": "top 50 upi banks volume value since inception",
}
eco_specs = [
    ("eco-upi-apps", "eco_upi_apps", "UPI app-wise transactions (CIT/B2C/B2B/Onus)",
     "Per-app monthly volumes and values split by transaction type, all TPAP apps on UPI.",
     NPCI_ECO_URL, "npci_eco_fetch.py"),
    ("eco-chargeback", "eco_chargeback", "UPI chargebacks by bank",
     "Chargebacks raised, representments and acceptances per bank code with CB ratio.",
     NPCI_ECO_URL, "npci_eco_fetch.py"),
    ("eco-top50-member", "eco_top50_member", "Top remitter/beneficiary bank performance",
     "Top-50 remitter and beneficiary banks monthly: total volume, approved %, business decline %, technical decline %.",
     NPCI_ECO_URL, "npci_eco_fetch.py"),
    ("eco-psp", "eco_psp", "UPI payer/payee PSP performance",
     "Payer- and payee-side PSP monthly performance: volume, approved %, BD %, TD %.",
     NPCI_ECO_URL, "npci_eco_fetch.py"),
    ("eco-mcc", "eco_mcc", "UPI merchant category (MCC) classification",
     "High/low transacting merchant categories with volume and value contribution.",
     NPCI_ECO_URL, "npci_eco_fetch.py"),
    ("eco-p2p-2m", "eco_p2p_p2m", "UPI P2P vs P2M split",
     "Monthly P2P and P2M volumes and values.",
     NPCI_ECO_URL, "npci_eco_fetch.py"),
    ("eco-statewise", "eco_statewise", "UPI state-wise transactions",
     "Per-state UPI volume/value with contribution %.",
     NPCI_ECO_URL, "npci_eco_fetch.py"),
    ("eco-top50-volval", "eco_top50_volval", "Top-50 UPI banks by volume & value",
     "Monthly top-50 UPI member banks by volume and value (inception to date).",
     NPCI_ECO_URL, "npci_eco_fetch.py"),
]
for cid, table, title, desc, src, script in eco_specs:
    if not con.execute("SELECT count(*) FROM information_schema.tables WHERE table_name=?", [table]).fetchone()[0]:
        continue
    cov = con.execute(f"SELECT min(month), max(month), count(*) FROM {table}").fetchone()
    if cov[0] is None:
        continue
    add(cid, title, desc, "table+chart", src,
        script, {"from": str(cov[0]), "to": str(cov[1]), "rows": cov[2]},
        extra={"api": f"/api/payments-hub?dataset=eco&tab={cid.replace('eco-','')}"})

for e in cat:
    e["aliases"] = ALIASES.get(e["id"], "")

dump("catalog.json", cat)
print(f"catalog: {len(cat)} datasets")
print("DONE")
