CoolFace
Apppublic

mrik8899/fintech-demo

sourceHugging Faceupdated 5mo agoView on Hugging Face
0likes
db.py1043 linesDownload Raw Back to app
1# app/db.py2import os3import time4import threading5from decimal import Decimal6from datetime import datetime, timedelta7from app.config import DB_PATH, DATABASE_URL8import psycopg29from psycopg2 import pool, extras10from urllib.parse import urlparse, parse_qs, urlencode, urlunparse11 12 13def _dsn_with_keepalives(dsn: str) -> str:14    """Append TCP keepalive params to prevent Neon from dropping idle SSL connections."""15    parsed = urlparse(dsn)16    params = parse_qs(parsed.query, keep_blank_values=True)17    params.setdefault("sslmode", ["require"])18    params["keepalives"] = ["1"]19    params["keepalives_idle"] = ["30"]20    params["keepalives_interval"] = ["10"]21    params["keepalives_count"] = ["5"]22    new_query = urlencode({k: v[0] for k, v in params.items()})23    return urlunparse(parsed._replace(query=new_query))24 25# ── Connection Management ─────────────────────────────────────────────────26 27_pool = None28_db_initialized = False29 30def _ensure_pool():31    global _pool32    if _pool:33        return34    if DATABASE_URL:35        print("DB: Connecting to PostgreSQL (Neon)...")36        last_err = None37        for attempt in range(5):38            try:39                _pool = pool.ThreadedConnectionPool(40                    minconn=2,41                    maxconn=10,42                    dsn=_dsn_with_keepalives(DATABASE_URL)43                )44                print("DB: Connection pool ready.")45                return46            except Exception as e:47                last_err = e48                wait = (attempt + 1) * 3  # 3s, 6s, 9s, 12s, 15s49                print(f"DB: Pool init attempt {attempt + 1}/5 failed "50                      f"({e}) — retrying in {wait}s...")51                time.sleep(wait)52        print(f"DB: All pool init attempts failed — {last_err}")53    else:54        print("DB: DATABASE_URL not set, falling back to SQLite.")55 56 57def get_connection():58    _ensure_pool()59    if _pool:60        conn = _pool.getconn()61        conn.autocommit = False62        # Validate connection — Neon drops idle SSL connections after ~5min63        try:64            cur = conn.cursor()65            cur.execute("SELECT 1")66            cur.close()                    # ← close cursor (fixes Bug 7 leak)67        except Exception:68            # Connection is dead — close it and get a fresh one from pool69            try:70                _pool.putconn(conn, close=True)71            except Exception:72                pass73            conn = _pool.getconn()74            conn.autocommit = False75        return conn76    else:77        import sqlite378        base_dir = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))79        db_full_path = os.path.join(base_dir, DB_PATH)80        conn = sqlite3.connect(db_full_path, timeout=30)81        conn.row_factory = sqlite3.Row82        conn.execute("PRAGMA journal_mode=WAL")83        conn.execute("PRAGMA busy_timeout=30000")84        return conn85    86def release_connection(conn):87    if _pool and conn:88        _pool.putconn(conn)89 90def close_pool():91    global _pool92    if _pool:93        _pool.closeall()94        _pool = None95 96# ── Helpers ───────────────────────────────────────────────────────────────97 98def _is_postgres(conn):99    return hasattr(conn, 'autocommit') and DATABASE_URL100 101def _fetch_all(conn, cursor):102    if _is_postgres(conn):103        cols = [desc[0] for desc in cursor.description]104        rows = []105        for row in cursor.fetchall():106            d = {}107            for i, col in enumerate(cols):108                val = row[i]109                if isinstance(val, Decimal):110                    val = float(val)111                d[col] = val112            rows.append(d)113        return rows114    else:115        return [dict(r) for r in cursor.fetchall()]116 117# ── Schema Migrations ────────────────────────────────────────────────────118 119def _ensure_day_index():120    global _db_initialized121    if _db_initialized:122        return123    conn = get_connection()124    cursor = conn.cursor()125    try:126        if _is_postgres(conn):127            cursor.execute("""128                DO $$ BEGIN129                    ALTER TABLE transactions ADD COLUMN IF NOT EXISTS day TEXT;130                EXCEPTION WHEN duplicate_column THEN NULL;131                END $$;132            """)133            cursor.execute("""134                CREATE INDEX IF NOT EXISTS idx_txn_merchant_day ON transactions(merchant_id, day);135            """)136        else:137            cursor.execute("PRAGMA table_info(transactions)")138            columns = [col[1] for col in cursor.fetchall()]139            if "day" not in columns:140                print("  DB: One-time migration — adding day column and index...")141                cursor.execute("ALTER TABLE transactions ADD COLUMN day TEXT")142                cursor.execute("UPDATE transactions SET day = substr(timestamp, 1, 10)")143                cursor.execute("CREATE INDEX idx_txn_merchant_day ON transactions(merchant_id, day)")144        conn.commit()145        _db_initialized = True146    finally:147        release_connection(conn)148 149# ── Query Functions ──────────────────────────────────────────────────────150 151def _execute_query(query, params=None, fetch=True):152    conn = get_connection()153    try:154        cursor = conn.cursor()155        cursor.execute(query, params)156        if fetch:157            res = _fetch_all(conn, cursor)158            return res159        return None160    finally:161        conn.commit()162        release_connection(conn)163 164 165def _safe_growth_pct(current, previous, max_pct=999.9):166    """Calculate growth %, capped to avoid absurd numbers from thin comparison periods.167    168    When previous period has very little data, raw (curr-prev)/prev produces169    billions of percent. Cap at ±999.9% — still signals "huge shift" without170    being absurd. Returns 0.0 when both periods are zero.171    """172    if previous <= 0:173        return max_pct if current > 0 else 0.0174    pct = (current - previous) / previous * 100175    return round(max(min(pct, max_pct), -max_pct), 1)176 177 178def _safe_share_change_pp(share_now, share_prev, prev_amount, max_pp=50.0):179    """Calculate share change in pp, capped and suppressed when no previous data.180    181    If previous period has zero amount, we can't meaningfully say share "shifted"182    — it's just new data appearing. Return 0 so the UI hides the arrow line.183    Otherwise cap at ±50pp to avoid absurd jumps.184    """185    if prev_amount <= 0:186        return 0.0187    pp = share_now - share_prev188    return round(max(min(pp, max_pp), -max_pp), 1)189 190 191def get_daily_aggregates(days=35, reference_date=None):192    _ensure_day_index()193    if reference_date:194        cutoff = (reference_date - timedelta(days=days)).strftime("%Y-%m-%d")195    else:196        cutoff = (datetime.now() - timedelta(days=days)).strftime("%Y-%m-%d")197    198    query = """199        SELECT200            merchant_id,201            day,202            COUNT(*) AS txn_count,203            SUM(CASE WHEN status = 'success' THEN 1 ELSE 0 END) AS success_count,204            SUM(CASE WHEN status = 'failed' THEN 1 ELSE 0 END) AS failed_count,205            SUM(CASE WHEN status = 'declined' THEN 1 ELSE 0 END) AS declined_count,206            ROUND(SUM(amount)::numeric, 2) AS revenue,207            ROUND(SUM(CASE WHEN status = 'success' THEN amount ELSE 0 END)::numeric, 2) AS success_revenue,208            ROUND(AVG(CASE WHEN status != 'success' THEN bank_latency_ms END), 0) AS avg_fail_bank_latency209        FROM transactions210        WHERE day >= %s211        GROUP BY merchant_id, day212        ORDER BY merchant_id, day213"""214    return _execute_query(query, (cutoff,))215 216 217def get_all_transactions():218    query = "SELECT id, merchant_id, amount, status, timestamp FROM transactions ORDER BY timestamp ASC LIMIT 10000"219    return _execute_query(query)220 221 222def get_all_merchants():223    query = """224        SELECT id, name, city, segment, category, pattern, status, contact_person, phone, created_at, 225        acquiring_bank, integrator_type, settlement_cycle, mdr_rate, base_daily_vol, base_avg_amount 226        FROM merchants227    """228    rows = _execute_query(query)229    return {r["id"]: r for r in rows}230 231 232def get_merchant_transactions(merchant_id: int):233    query = """234        SELECT id, merchant_id, amount, status, timestamp, payment_method, failure_reason, response_code, 235        gateway_latency_ms, bank_latency_ms, settlement_status 236        FROM transactions WHERE merchant_id = %s ORDER BY timestamp DESC LIMIT 500237"""238    return _execute_query(query, (merchant_id,))239 240 241def get_failure_breakdown(merchant_id: int, days: int = 7, reference_date=None):242    if reference_date:243        cutoff = (reference_date - timedelta(days=days)).strftime("%Y-%m-%d")244    else:245        cutoff = (datetime.now() - timedelta(days=days)).strftime("%Y-%m-%d")246    247    p = "%s"248    query = f"""249        SELECT failure_reason, COUNT(*) as cnt, ROUND(AVG(bank_latency_ms), 0) as avg_bank_lat 250        FROM transactions WHERE merchant_id = {p} AND status != 'success' AND day >= {p} 251        GROUP BY failure_reason ORDER BY cnt DESC LIMIT 5252    """253    return _execute_query(query, (merchant_id, cutoff))254 255 256def get_payment_breakdown(merchant_id: int, days: int = 7, reference_date=None):257    if reference_date:258        cutoff = (reference_date - timedelta(days=days)).strftime("%Y-%m-%d")259    else:260        cutoff = (datetime.now() - timedelta(days=days)).strftime("%Y-%m-%d")261    262    p = "%s"263    query = f"""264        SELECT payment_method, COUNT(*) as cnt, ROUND(SUM(amount), 0) as total 265        FROM transactions WHERE merchant_id = {p} AND day >= {p} 266        GROUP BY payment_method ORDER BY cnt DESC267    """268    return _execute_query(query, (merchant_id, cutoff))269 270 271def get_settlement_summary(reference_date=None):272    if reference_date:273        cutoff = (reference_date - timedelta(days=3)).strftime("%Y-%m-%d")274    else:275        cutoff = (datetime.now() - timedelta(days=3)).strftime("%Y-%m-%d")276    277    p = "%s"278    query = f"""279        SELECT settlement_status, COUNT(*) as cnt, ROUND(SUM(amount), 0) as total 280        FROM transactions WHERE day >= {p} AND status = 'success' 281        GROUP BY settlement_status282    """283    return _execute_query(query, (cutoff,))284 285 286def get_settlement_health(reference_date=None, days=7):287    if reference_date:288        cutoff = (reference_date - timedelta(days=days)).strftime("%Y-%m-%d")289    else:290        cutoff = (datetime.now() - timedelta(days=days)).strftime("%Y-%m-%d")291    292    p = "%s"293    conn = get_connection()294    try:295        cursor = conn.cursor()296 297        cursor.execute(f"""298            SELECT settlement_status, COUNT(*) as cnt, ROUND(SUM(amount), 0) as total299            FROM transactions WHERE day >= {p} AND status = 'success'300            GROUP BY settlement_status301        """, (cutoff,))302        status_breakdown = _fetch_all(conn, cursor)303 304        cursor.execute(f"""305            SELECT m.acquiring_bank as bank, t.settlement_status,306                   COUNT(*) as cnt, ROUND(SUM(t.amount), 0) as total307            FROM transactions t308            JOIN merchants m ON t.merchant_id = m.id309            WHERE t.day >= {p} AND t.status = 'success'310            GROUP BY m.acquiring_bank, t.settlement_status311            ORDER BY m.acquiring_bank, t.settlement_status312        """, (cutoff,))313        bank_rows = _fetch_all(conn, cursor)314        by_bank = {}315        for r in bank_rows:316            bank = r["bank"]317            if bank not in by_bank: by_bank[bank] = []318            by_bank[bank].append({"status": r["settlement_status"], "count": r["cnt"], "amount": r["total"]})319 320        cursor.execute(f"""321            SELECT m.id as merchant_id, m.name as merchant_name, m.acquiring_bank, m.settlement_cycle,322                   COUNT(*) as pending_count, ROUND(SUM(t.amount), 0) as pending_amount323            FROM transactions t JOIN merchants m ON t.merchant_id = m.id324            WHERE t.day >= {p} AND t.status = 'success' AND t.settlement_status = 'pending'325            GROUP BY t.merchant_id, m.id, m.name, m.acquiring_bank, m.settlement_cycle326            ORDER BY pending_amount DESC LIMIT 10327        """, (cutoff,))328        top_pending = _fetch_all(conn, cursor)329 330        return {"period": f"last {days} days", "status_breakdown": status_breakdown, "by_bank": by_bank, "top_pending_merchants": top_pending}331    finally:332        conn.commit()333        release_connection(conn)334 335 336def get_mdr_revenue(reference_date=None, days=7, bank=None, category=None):337    if reference_date:338        cutoff = (reference_date - timedelta(days=days)).strftime("%Y-%m-%d")339    else:340        cutoff = (datetime.now() - timedelta(days=days)).strftime("%Y-%m-%d")341 342    join_sql = "JOIN merchants m ON t.merchant_id = m.id"343    conditions = ["t.day >= %s"]  # Always %s — PostgreSQL only344    params: list = [cutoff]345 346    if bank:347        conditions.append("m.acquiring_bank = %s")348        params.append(bank)349    if category:350        conditions.append("m.category = %s")351        params.append(category)352 353    where = " AND ".join(conditions)354 355    conn = get_connection()356    try:357        cursor = conn.cursor()358 359        cursor.execute(f"""360            SELECT COUNT(*) as success_count, ROUND(SUM(t.amount), 0) as success_amount,361                   ROUND(SUM(t.amount * m.mdr_rate / 100.0), 0) as mdr_earned362            FROM transactions t {join_sql} WHERE {where} AND t.status = 'success'363        """, params)364        earned = dict(cursor.fetchone())365 366        cursor.execute(f"""367            SELECT COUNT(*) as failed_count, ROUND(SUM(t.amount), 0) as failed_amount,368                   ROUND(SUM(t.amount * m.mdr_rate / 100.0), 0) as mdr_lost369            FROM transactions t {join_sql} WHERE {where} AND t.status != 'success'370        """, params)371        lost = dict(cursor.fetchone())372 373        by_bank = None374        if not bank:375            cursor.execute(f"""376                SELECT m.acquiring_bank as bank, COUNT(*) as success_count,377                       ROUND(SUM(t.amount), 0) as success_amount,378                       ROUND(SUM(t.amount * m.mdr_rate / 100.0), 0) as mdr_earned379                FROM transactions t {join_sql} WHERE {where} AND t.status = 'success'380                GROUP BY m.acquiring_bank ORDER BY mdr_earned DESC381            """, params)382            bank_earned = {r["bank"]: dict(r) for r in _fetch_all(conn, cursor)}383 384            cursor.execute(f"""385                SELECT m.acquiring_bank as bank, COUNT(*) as failed_count,386                       ROUND(SUM(t.amount), 0) as failed_amount,387                       ROUND(SUM(t.amount * m.mdr_rate / 100.0), 0) as mdr_lost388                FROM transactions t {join_sql} WHERE {where} AND t.status != 'success'389                GROUP BY m.acquiring_bank ORDER BY mdr_lost DESC390            """, params)391            for r in _fetch_all(conn, cursor):392                b = r["bank"]393                if b in bank_earned:394                    bank_earned[b].update(dict(r))395                else:396                    bank_earned[b] = dict(r)397            by_bank = bank_earned398 399        cursor.execute(f"""400            SELECT m.id as merchant_id, m.name as merchant_name,401                   m.acquiring_bank, m.category, m.mdr_rate,402                   COUNT(*) as failed_count,403                   ROUND(SUM(t.amount), 0) as failed_amount,404                   ROUND(SUM(t.amount * m.mdr_rate / 100.0), 0) as mdr_lost405            FROM transactions t {join_sql} WHERE {where} AND t.status != 'success'406            GROUP BY t.merchant_id, m.id, m.name, m.acquiring_bank, m.category, m.mdr_rate407            ORDER BY mdr_lost DESC LIMIT 10408        """, params)409        top_loss = _fetch_all(conn, cursor)410 411        mdr_earned = earned["mdr_earned"] or 0412        mdr_lost = lost["mdr_lost"] or 0413        return {414            "period": f"last {days} days",415            "filter": {"bank": bank, "category": category},416            "mdr_earned": mdr_earned,417            "mdr_lost": mdr_lost,418            "mdr_loss_rate": round(mdr_lost / max(mdr_earned + mdr_lost, 1), 3),419            "success_count": earned["success_count"] or 0,420            "success_amount": earned["success_amount"] or 0,421            "failed_count": lost["failed_count"] or 0,422            "failed_amount": lost["failed_amount"] or 0,423            "by_bank": by_bank,424            "top_mdr_loss_merchants": top_loss,425        }426    finally:427        conn.commit()428        release_connection(conn)429 430def get_raast_metrics(reference_date=None, days=7, bank=None):431    if reference_date:432        cutoff = (reference_date - timedelta(days=days)).strftime("%Y-%m-%d")433    else:434        cutoff = (datetime.now() - timedelta(days=days)).strftime("%Y-%m-%d")435 436    join_sql = "JOIN merchants m ON t.merchant_id = m.id"437    conditions = ["t.day >= %s", "t.payment_method = 'raast'"]438    params: list = [cutoff]439    if bank:440        conditions.append("m.acquiring_bank = %s")441        params.append(bank)442    where = " AND ".join(conditions)443    p = "%s"444 445    conn = get_connection()446    try:447        cursor = conn.cursor()448        cursor.execute(f"""449            SELECT COUNT(*) as total_count,450                   SUM(CASE WHEN t.status = 'success' THEN 1 ELSE 0 END) as success_count,451                   SUM(CASE WHEN t.status != 'success' THEN 1 ELSE 0 END) as failed_count,452                   ROUND(SUM(t.amount), 0) as total_amount,453                   ROUND(SUM(CASE WHEN t.status = 'success' THEN t.amount ELSE 0 END), 0) as success_amount,454                   ROUND(AVG(CASE WHEN t.status = 'success' THEN t.bank_latency_ms END), 0) as avg_success_latency,455                   ROUND(AVG(CASE WHEN t.status != 'success' THEN t.bank_latency_ms END), 0) as avg_fail_latency456            FROM transactions t {join_sql} WHERE {where}457        """, params)458        overall = dict(cursor.fetchone())459 460        cursor.execute(f"""461            SELECT t.day, COUNT(*) as total_count,462                   SUM(CASE WHEN t.status = 'success' THEN 1 ELSE 0 END) as success_count,463                   ROUND(SUM(t.amount), 0) as total_amount,464                   ROUND(SUM(CASE WHEN t.status = 'success' THEN t.amount ELSE 0 END), 0) as success_amount465            FROM transactions t {join_sql} WHERE {where} GROUP BY t.day ORDER BY t.day466        """, params)467        daily_trend = _fetch_all(conn, cursor)468 469        cursor.execute(f"""470            SELECT t.failure_reason, t.response_code, COUNT(*) as cnt, ROUND(SUM(t.amount), 0) as failed_amount471            FROM transactions t {join_sql} WHERE {where} AND t.status != 'success' AND t.failure_reason IS NOT NULL472            GROUP BY t.failure_reason, t.response_code ORDER BY cnt DESC473        """, params)474        failure_breakdown = _fetch_all(conn, cursor)475 476        by_bank = None477        if not bank:478            cursor.execute(f"""479                SELECT m.acquiring_bank as bank, COUNT(*) as total_count,480                       SUM(CASE WHEN t.status = 'success' THEN 1 ELSE 0 END) as success_count,481                       ROUND(SUM(t.amount), 0) as total_amount, ROUND(AVG(t.bank_latency_ms), 0) as avg_latency482                FROM transactions t {join_sql} WHERE {where} GROUP BY m.acquiring_bank ORDER BY total_amount DESC483            """, params)484            by_bank = _fetch_all(conn, cursor)485            for b in by_bank: b["success_rate"] = round(b["success_count"] / max(b["total_count"], 1), 3)486 487        cursor.execute(f"""488            SELECT t.settlement_status, COUNT(*) as cnt, ROUND(SUM(t.amount), 0) as total489            FROM transactions t {join_sql} WHERE {where} AND t.status = 'success' GROUP BY t.settlement_status490        """, params)491        settlement_status = _fetch_all(conn, cursor)492 493        cursor.execute(f"""494            SELECT m.id as merchant_id, m.name as merchant_name, m.acquiring_bank, m.category,495                   COUNT(*) as raast_count, ROUND(SUM(t.amount), 0) as raast_amount,496                   SUM(CASE WHEN t.status = 'success' THEN 1 ELSE 0 END) as success_count497            FROM transactions t {join_sql} WHERE {where}498            GROUP BY t.merchant_id, m.id, m.name, m.acquiring_bank, m.category ORDER BY raast_amount DESC LIMIT 10499        """, params)500        top_merchants = _fetch_all(conn, cursor)501        for m in top_merchants: m["success_rate"] = round(m["success_count"] / max(m["raast_count"], 1), 3)502 503        sc = overall["success_count"] or 0504        tc = overall["total_count"] or 0505        return {506            "period": f"last {days} days", "filter": {"bank": bank},507            "total_count": tc, "success_count": sc, "failed_count": overall["failed_count"] or 0,508            "success_rate": round(sc / max(tc, 1), 3), "failure_rate": round(1 - sc / max(tc, 1), 3),509            "total_amount": overall["total_amount"] or 0, "success_amount": overall["success_amount"] or 0,510            "avg_success_latency": overall["avg_success_latency"] or 0, "avg_fail_latency": overall["avg_fail_latency"] or 0,511            "daily_trend": daily_trend, "failure_breakdown": failure_breakdown, "by_bank": by_bank,512            "settlement_status": settlement_status, "top_merchants": top_merchants,513        }514    finally:515        conn.commit()516        release_connection(conn)517 518 519def get_payment_mix_trend(reference_date=None, days=30, group_by=None):520    _VALID_GROUP_BY = {"category", "city", "bank"}521    if group_by not in _VALID_GROUP_BY: group_by = None522    if reference_date: ref = reference_date523    else: ref = datetime.now()524 525    period_start = (ref - timedelta(days=days)).strftime("%Y-%m-%d")526    prev_start = (ref - timedelta(days=days * 2)).strftime("%Y-%m-%d")527    p = "%s"528 529    need_join = group_by in ("category", "city", "bank")530    join_sql = "JOIN merchants m ON t.merchant_id = m.id" if need_join else ""531    if group_by == "category": group_cols = "m.category, t.payment_method"; _group_key_col = "category"532    elif group_by == "city": group_cols = "m.city, t.payment_method"; _group_key_col = "city"533    elif group_by == "bank": group_cols = "m.acquiring_bank, t.payment_method"; _group_key_col = "acquiring_bank"534    else: group_cols = "t.payment_method"; _group_key_col = None535 536    conn = get_connection()537    try:538        cursor = conn.cursor()539        cursor.execute(f"SELECT {group_cols} as method_key, COUNT(*) as txn_count, ROUND(SUM(t.amount), 0) as total_amount FROM transactions t {join_sql} WHERE t.day >= {p} AND t.status = 'success' GROUP BY {group_cols} ORDER BY total_amount DESC", (period_start,))540        current_rows = _fetch_all(conn, cursor)541        cursor.execute(f"SELECT {group_cols} as method_key, COUNT(*) as txn_count, ROUND(SUM(t.amount), 0) as total_amount FROM transactions t {join_sql} WHERE t.day >= {p} AND t.day < {p} AND t.status = 'success' GROUP BY {group_cols} ORDER BY total_amount DESC", (prev_start, period_start))542        prev_rows = _fetch_all(conn, cursor)543    finally:544        conn.commit()545        release_connection(conn)546 547    total_current = sum(r["total_amount"] for r in current_rows)548    total_previous = sum(r["total_amount"] for r in prev_rows)549    prev_map = {}550    for r in prev_rows:551        pkey = f"{r[_group_key_col]} · {r['method_key']}" if _group_key_col else r["method_key"]552        prev_map[pkey] = r553    comparison, seen = [], set()554    for r in current_rows:555        key = f"{r[_group_key_col]} · {r['method_key']}" if _group_key_col else r["method_key"]556        seen.add(key)557        prev = prev_map.get(key, {"total_amount": 0, "txn_count": 0})558        share_now = (r["total_amount"] / max(total_current, 1)) * 100559        share_prev = (prev["total_amount"] / max(total_previous, 1)) * 100560        comparison.append({561            "method": key,562            "amount": r["total_amount"],563            "txn_count": r["txn_count"],564            "share_pct": round(share_now, 1),565            "prev_amount": prev["total_amount"],566            "prev_share_pct": round(share_prev, 1),567            "share_change_pp": _safe_share_change_pp(share_now, share_prev, prev["total_amount"]),568            "growth_pct": _safe_growth_pct(r["total_amount"], prev["total_amount"])569        })570    _prev_key_col = _group_key_col  # Same column mapping for previous period571    for r in prev_rows:572        pkey = f"{r[_prev_key_col]} · {r['method_key']}" if _prev_key_col else r["method_key"]573        if pkey not in seen:574            share_prev = (r["total_amount"] / max(total_previous, 1)) * 100575            comparison.append({"method": pkey, "amount": 0, "txn_count": 0, "share_pct": 0, "prev_amount": r["total_amount"], "prev_share_pct": round(share_prev, 1), "share_change_pp": round(-share_prev, 1), "growth_pct": -100.0})576    comparison.sort(key=lambda x: x["amount"], reverse=True)577    return {578        "period_days": days, "group_by": group_by,579        "total_gmv_current": total_current, "total_gmv_previous": total_previous,580        "overall_growth_pct": _safe_growth_pct(total_current, total_previous),581        "mix_shift": comparison[:25]582    }583 584 585def get_gmv_growth(reference_date=None, days=30, group_by=None):586    _VALID_GROUP_BY = {"category", "city", "bank", "segment"}587    if group_by not in _VALID_GROUP_BY: group_by = None588    if reference_date: ref = reference_date589    else: ref = datetime.now()590    period_start = (ref - timedelta(days=days)).strftime("%Y-%m-%d")591    prev_start = (ref - timedelta(days=days * 2)).strftime("%Y-%m-%d")592    p = "%s"593 594    if group_by == "city": group_col = "m.city"595    elif group_by == "bank": group_col = "m.acquiring_bank"596    elif group_by == "segment": group_col = "m.segment"597    else: group_col = "m.category"; group_by = "category"598 599    conn = get_connection()600    try:601        cursor = conn.cursor()602        cursor.execute(f"SELECT {group_col} as group_name, COUNT(*) as txn_count, ROUND(SUM(t.amount), 0) as total_amount, ROUND(SUM(t.amount) / COUNT(DISTINCT t.merchant_id), 0) as gmv_per_merchant FROM transactions t JOIN merchants m ON t.merchant_id = m.id WHERE t.day >= {p} AND t.status = 'success' GROUP BY {group_col} ORDER BY total_amount DESC", (period_start,))603        current_rows = _fetch_all(conn, cursor)604        cursor.execute(f"SELECT {group_col} as group_name, COUNT(*) as txn_count, ROUND(SUM(t.amount), 0) as total_amount FROM transactions t JOIN merchants m ON t.merchant_id = m.id WHERE t.day >= {p} AND t.day < {p} AND t.status = 'success' GROUP BY {group_col}", (prev_start, period_start))605        prev_rows = _fetch_all(conn, cursor)606        cursor.execute("SELECT COUNT(DISTINCT merchant_id) as active_merchants FROM transactions WHERE day >= %s AND status = 'success'", (period_start,))607        res = cursor.fetchone()608        active_merchants = res[0] if _is_postgres(conn) else dict(res)["active_merchants"]609    finally:610        conn.commit()611        release_connection(conn)612 613    total_current = sum(r["total_amount"] for r in current_rows)614    total_previous = sum(r["total_amount"] for r in prev_rows)615    prev_map = {r["group_name"]: r for r in prev_rows}616    growth_data = []617    for r in current_rows:618        name = r["group_name"]; prev = prev_map.get(name, {"total_amount": 0, "txn_count": 0})619        growth_data.append({620            "group": name,621            "amount": r["total_amount"],622            "txn_count": r["txn_count"],623            "gmv_per_merchant": r["gmv_per_merchant"],624            "share_pct": round((r["total_amount"] / max(total_current, 1)) * 100, 1),625            "prev_amount": prev["total_amount"],626            "prev_txn_count": prev["txn_count"],627            "growth_pct": _safe_growth_pct(r["total_amount"], prev["total_amount"]),628            "txn_growth_pct": round((r["txn_count"] - prev["txn_count"]) / max(prev["txn_count"], 1) * 100, 1)629        })630    growth_data.sort(key=lambda x: x["growth_pct"], reverse=True)631    return {632        "period_days": days, "group_by": group_by,633        "total_gmv_current": total_current, "total_gmv_previous": total_previous,634        "overall_growth_pct": _safe_growth_pct(total_current, total_previous),635        "active_merchants": active_merchants, "growth_breakdown": growth_data636    }637 638 639def get_bank_market_share(reference_date=None, days=30):640    if reference_date: ref = reference_date641    else: ref = datetime.now()642    period_start = (ref - timedelta(days=days)).strftime("%Y-%m-%d")643    prev_start = (ref - timedelta(days=days * 2)).strftime("%Y-%m-%d")644    p = "%s"645 646    conn = get_connection()647    try:648        cursor = conn.cursor()649        cursor.execute(f"SELECT m.acquiring_bank as bank, COUNT(*) as txn_count, ROUND(SUM(t.amount), 0) as total_amount, SUM(CASE WHEN t.status = 'success' THEN 1 ELSE 0 END) as success_count, ROUND(AVG(t.bank_latency_ms), 0) as avg_latency FROM transactions t JOIN merchants m ON t.merchant_id = m.id WHERE t.day >= {p} GROUP BY m.acquiring_bank ORDER BY total_amount DESC", (period_start,))650        current_rows = _fetch_all(conn, cursor)651        cursor.execute(f"SELECT m.acquiring_bank as bank, COUNT(*) as txn_count, ROUND(SUM(t.amount), 0) as total_amount FROM transactions t JOIN merchants m ON t.merchant_id = m.id WHERE t.day >= {p} AND t.day < {p} GROUP BY m.acquiring_bank", (prev_start, period_start))652        prev_rows = _fetch_all(conn, cursor)653    finally:654        conn.commit()655        release_connection(conn)656 657    total_current = sum(r["total_amount"] for r in current_rows)658    total_previous = sum(r["total_amount"] for r in prev_rows)659    prev_map = {r["bank"]: r for r in prev_rows}660    banks = []661    for r in current_rows:662        b = r["bank"]; prev = prev_map.get(b, {"total_amount": 0, "txn_count": 0})663        share_now = (r["total_amount"] / max(total_current, 1)) * 100664        share_prev = (prev["total_amount"] / max(total_previous, 1)) * 100665        banks.append({666            "bank": b,667            "amount": r["total_amount"],668            "txn_count": r["txn_count"],669            "success_rate": round(r["success_count"] / max(r["txn_count"], 1), 3),670            "avg_latency_ms": r["avg_latency"],671            "share_pct": round(share_now, 1),672            "prev_share_pct": round(share_prev, 1),673            "share_change_pp": _safe_share_change_pp(share_now, share_prev, prev["total_amount"]),674            "growth_pct": _safe_growth_pct(r["total_amount"], prev["total_amount"]),675            "rank": 0676        })677    for r in prev_rows:678        if r["bank"] not in {b["bank"] for b in banks}:679            share_prev = (r["total_amount"] / max(total_previous, 1)) * 100680            banks.append({"bank": r["bank"], "amount": 0, "txn_count": 0, "success_rate": 0, "avg_latency_ms": 0, "share_pct": 0, "prev_share_pct": round(share_prev, 1), "share_change_pp": round(-share_prev, 1), "growth_pct": -100.0, "rank": 0})681    banks.sort(key=lambda x: x["amount"], reverse=True)682    for i, b in enumerate(banks): b["rank"] = i + 1683    return {684        "period_days": days,685        "total_gmv_current": total_current, "total_gmv_previous": total_previous,686        "overall_growth_pct": _safe_growth_pct(total_current, total_previous),687        "banks": banks688    }689 690 691def get_segment_analysis(reference_date=None, days=30):692    if reference_date: ref = reference_date693    else: ref = datetime.now()694    period_start = (ref - timedelta(days=days)).strftime("%Y-%m-%d")695    prev_start = (ref - timedelta(days=days * 2)).strftime("%Y-%m-%d")696    p = "%s"697 698    conn = get_connection()699    try:700        cursor = conn.cursor()701        cursor.execute(f"SELECT m.segment, COUNT(DISTINCT m.id) as merchant_count, COUNT(*) as txn_count, ROUND(SUM(t.amount), 0) as total_amount, ROUND(SUM(t.amount) / COUNT(DISTINCT m.id), 0) as gmv_per_merchant, ROUND(AVG(t.amount), 0) as avg_ticket FROM transactions t JOIN merchants m ON t.merchant_id = m.id WHERE t.day >= {p} AND t.status = 'success' GROUP BY m.segment ORDER BY total_amount DESC", (period_start,))702        current_rows = _fetch_all(conn, cursor)703        cursor.execute(f"SELECT m.segment, COUNT(DISTINCT m.id) as merchant_count, ROUND(SUM(t.amount), 0) as total_amount FROM transactions t JOIN merchants m ON t.merchant_id = m.id WHERE t.day >= {p} AND t.day < {p} AND t.status = 'success' GROUP BY m.segment", (prev_start, period_start))704        prev_rows = _fetch_all(conn, cursor)705    finally:706        conn.commit()707        release_connection(conn)708 709    total_current = sum(r["total_amount"] for r in current_rows)710    total_previous = sum(r["total_amount"] for r in prev_rows)711    prev_map = {r["segment"]: r for r in prev_rows}712    segments = []713    for r in current_rows:714        s = r["segment"]; prev = prev_map.get(s, {"total_amount": 0, "merchant_count": 0})715        segments.append({716            "segment": s,717            "merchant_count": r["merchant_count"],718            "txn_count": r["txn_count"],719            "total_gmv": r["total_amount"],720            "gmv_per_merchant": r["gmv_per_merchant"],721            "avg_ticket": r["avg_ticket"],722            "share_pct": round((r["total_amount"] / max(total_current, 1)) * 100, 1),723            "prev_gmv": prev["total_amount"],724            "prev_merchants": prev["merchant_count"],725            "gmv_growth_pct": _safe_growth_pct(r["total_amount"], prev["total_amount"])726        })727    segments.sort(key=lambda x: x["gmv_growth_pct"], reverse=True)728    return {729        "period_days": days,730        "total_gmv_current": total_current, "total_gmv_previous": total_previous,731        "overall_growth_pct": _safe_growth_pct(total_current, total_previous),732        "segments": segments733    }734 735 736def create_industry_benchmarks_table():737    conn = get_connection()738    try:739        cursor = conn.cursor()740        if _is_postgres(conn):741            cursor.execute("""742                CREATE TABLE IF NOT EXISTS industry_benchmarks (743                    period TEXT NOT NULL, metric TEXT NOT NULL, value REAL NOT NULL, label TEXT NOT NULL,744                    source TEXT NOT NULL DEFAULT 'SBP Payment Systems Statistics',745                    updated_at TEXT NOT NULL DEFAULT (now() at time zone 'utc'),746                    PRIMARY KEY (period, metric)747                )748            """)749        else:750            cursor.execute("""CREATE TABLE IF NOT EXISTS industry_benchmarks (751                period TEXT NOT NULL, metric TEXT NOT NULL, value REAL NOT NULL, label TEXT NOT NULL,752                source TEXT NOT NULL DEFAULT 'SBP Payment Systems Statistics',753                updated_at TEXT NOT NULL DEFAULT (datetime('now')),754                PRIMARY KEY (period, metric)755            )""")756        conn.commit()757 758        cursor.execute("SELECT COUNT(*) as cnt FROM industry_benchmarks")759        res = cursor.fetchone()760        count = res[0] if _is_postgres(conn) else dict(res)["cnt"]761        if count == 0:762            seed_industry_benchmarks(conn, cursor)763    finally:764        release_connection(conn)765 766 767def seed_industry_benchmarks(conn, cursor=None):768    if not cursor: cursor = conn.cursor()769    rows = [770        ("2025-Q4", "raast_volume_share_pct", 12.0, "Raast share of digital payments volume"),771        ("2025-Q4", "card_volume_share_pct", 45.0, "Card (debit+credit) share of digital payments volume"),772        ("2025-Q4", "wallet_volume_share_pct", 20.0, "Mobile wallet (non-Raast) share"),773        ("2025-Q4", "bank_transfer_share_pct", 18.0, "IBFT/FT share of digital payments volume"),774        ("2025-Q4", "qr_volume_share_pct", 5.0, "QR code share of digital payments volume"),775        ("2025-Q4", "total_digital_gmv_trillion_pkr", 35.0, "Total digital payments GMV (PKR trillion)"),776        ("2025-Q4", "raast_gmv_trillion_pkr", 4.2, "Raast GMV (PKR trillion)"),777        ("2025-Q4", "card_gmv_trillion_pkr", 15.75, "Card GMV (PKR trillion)"),778        ("2025-Q4", "wallet_gmv_trillion_pkr", 7.0, "Mobile wallet GMV (PKR trillion)"),779        ("2025-Q4", "digital_txn_count_millions", 2800, "Total digital payment transactions (millions)"),780        ("2025-Q4", "raast_success_rate_pct", 94.5, "Raast transaction success rate"),781        ("2025-Q4", "card_success_rate_pct", 88.0, "Card transaction success rate"),782        ("2025-Q4", "avg_ticket_size_pkr", 12500, "Average digital payment ticket size (PKR)"),783        ("2025-Q4", "active_merchants_millions", 0.35, "Active digital merchants (millions)"),784    ]785    p = "%s"786    for row in rows:787        cursor.execute(f"INSERT INTO industry_benchmarks (period, metric, value, label) VALUES ({p}, {p}, {p}, {p}) ON CONFLICT DO NOTHING", row)788    conn.commit()789    print(f"  DB: Seeded {len(rows)} industry benchmarks")790 791 792def get_industry_benchmarks(period=None):793    p = "%s"794    if period:795        query = f"SELECT period, metric, value, label, source, updated_at FROM industry_benchmarks WHERE period = {p} ORDER BY metric"796        return _execute_query(query, (period,))797    else:798        return _execute_query("SELECT period, metric, value, label, source, updated_at FROM industry_benchmarks ORDER BY period DESC, metric")799 800 801 802def get_failure_summary(reference_date=None, days=7, bank=None, category=None, failure_reason=None):803    """Failure reason breakdown by reason, with per-bank split when no bank filter."""804    if reference_date:805        cutoff = (reference_date - timedelta(days=days)).strftime("%Y-%m-%d")806    else:807        cutoff = (datetime.now() - timedelta(days=days)).strftime("%Y-%m-%d")808 809    need_join = bank is not None or category is not None810    join_sql = "JOIN merchants m ON t.merchant_id = m.id" if need_join else ""811    p = "%s"  # Always use %s — we only use PostgreSQL now812    conditions = ["t.status != 'success'", f"t.day >= {p}", "t.failure_reason IS NOT NULL"]813    params: list = [cutoff]814 815    if bank:816        conditions.append(f"m.acquiring_bank = {p}")817        params.append(bank)818    if category:819        conditions.append(f"m.category = {p}")820        params.append(category)821    if failure_reason:822        conditions.append(f"t.failure_reason = {p}")823        params.append(failure_reason)824 825    where = " AND ".join(conditions)826 827    conn = get_connection()828    try:829        cursor = conn.cursor()830        cursor.execute(f"""831            SELECT t.failure_reason, COUNT(*) as cnt,832                   ROUND(SUM(t.amount), 0) as total_amount,833                   ROUND(AVG(t.bank_latency_ms), 0) as avg_bank_lat834            FROM transactions t {join_sql}835            WHERE {where}836            GROUP BY t.failure_reason837            ORDER BY cnt DESC838        """, params)839        rows = _fetch_all(conn, cursor)840 841        by_bank = None842        if not bank:843            cursor.execute(f"""844                SELECT m.acquiring_bank as bank, t.failure_reason, COUNT(*) as cnt845                FROM transactions t846                JOIN merchants m ON t.merchant_id = m.id847                WHERE t.status != 'success' AND t.day >= %s AND t.failure_reason IS NOT NULL848                GROUP BY m.acquiring_bank, t.failure_reason849                ORDER BY m.acquiring_bank, cnt DESC850            """, (cutoff,))851            bank_rows = _fetch_all(conn, cursor)852            by_bank = {}853            for r in bank_rows:854                b = r["bank"]855                if b not in by_bank:856                    by_bank[b] = []857                by_bank[b].append({"reason": r["failure_reason"], "count": r["cnt"]})858 859        return {860            "period": f"last {days} days",861            "filter": {"bank": bank, "category": category},862            "total_failures": sum(r["cnt"] for r in rows),863            "breakdown": rows,864            "by_bank": by_bank,865        }866    finally:867        conn.commit()868        release_connection(conn)869 870 871def _ensure_resolved_table():872    conn = get_connection()873    try:874        cursor = conn.cursor()875        if _is_postgres(conn):876            cursor.execute("""877                CREATE TABLE IF NOT EXISTS resolved_merchants (878                    merchant_id INTEGER PRIMARY KEY,879                    resolved_at TEXT NOT NULL DEFAULT (now() at time zone 'utc'),880                    note TEXT DEFAULT '',881                    resolved_by TEXT DEFAULT 'dashboard',882                    outcome TEXT DEFAULT 'unknown',883                    issue_type TEXT DEFAULT '',884                    daily_loss_at_resolution REAL DEFAULT 0,885                    days_to_resolve INTEGER DEFAULT 0886                )887            """)888            # Safe migration for existing table889            cursor.execute("""890                DO $$ BEGIN891                    ALTER TABLE resolved_merchants892                    ADD COLUMN IF NOT EXISTS outcome TEXT DEFAULT 'unknown';893                    ALTER TABLE resolved_merchants894                    ADD COLUMN IF NOT EXISTS issue_type TEXT DEFAULT '';895                    ALTER TABLE resolved_merchants896                    ADD COLUMN IF NOT EXISTS daily_loss_at_resolution REAL DEFAULT 0;897                    ALTER TABLE resolved_merchants898                    ADD COLUMN IF NOT EXISTS days_to_resolve INTEGER DEFAULT 0;899                EXCEPTION WHEN duplicate_column THEN NULL;900                END $$;901            """)902        conn.commit()903    finally:904        release_connection(conn)905 906def get_resolved_merchant_ids():907    _ensure_resolved_table()908    rows = _execute_query("SELECT merchant_id FROM resolved_merchants")909    return set(r["merchant_id"] for r in rows)910 911def resolve_merchant(912    merchant_id: int,913    note: str = "",914    resolved_by: str = "dashboard",915    outcome: str = "unknown",916    issue_type: str = "",917    daily_loss: float = 0,918):919    _ensure_resolved_table()920    p = "%s"921    if _pool:922        query = f"""923            INSERT INTO resolved_merchants924            (merchant_id, resolved_at, note, resolved_by,925             outcome, issue_type, daily_loss_at_resolution)926            VALUES ({p}, now() at time zone 'utc', {p}, {p}, {p}, {p}, {p})927            ON CONFLICT (merchant_id) DO UPDATE SET928                resolved_at = now() at time zone 'utc',929                note = EXCLUDED.note,930                resolved_by = EXCLUDED.resolved_by,931                outcome = EXCLUDED.outcome,932                issue_type = EXCLUDED.issue_type,933                daily_loss_at_resolution = EXCLUDED.daily_loss_at_resolution934        """935    else:936        query = f"""937            INSERT OR REPLACE INTO resolved_merchants938            (merchant_id, resolved_at, note, resolved_by,939             outcome, issue_type, daily_loss_at_resolution)940            VALUES ({p}, datetime('now'), {p}, {p}, {p}, {p}, {p})941        """942    _execute_query(943        query,944        (merchant_id, note or "", resolved_by, outcome, issue_type, daily_loss),945        fetch=False946    )947 948def unresolve_merchant(merchant_id: int):949    _ensure_resolved_table()950    p = "%s"951    _execute_query(f"DELETE FROM resolved_merchants WHERE merchant_id = {p}", (merchant_id,), fetch=False)952 953def get_resolved_log(limit: int = 50):954    _ensure_resolved_table()955    p = "%s"956    return _execute_query(f"SELECT merchant_id, resolved_at, note, resolved_by FROM resolved_merchants ORDER BY resolved_at DESC LIMIT {p}", (limit,))957 958 959def get_resolution_outcomes():960    """Aggregate resolution outcomes for feedback analysis."""961    _ensure_resolved_table()962    p = "%s"963    conn = get_connection()964    try:965        cursor = conn.cursor()966        cursor.execute("""967            SELECT968                outcome,969                issue_type,970                COUNT(*) as count,971                ROUND(AVG(daily_loss_at_resolution), 0) as avg_daily_loss,972                ROUND(AVG(days_to_resolve), 1) as avg_days_to_resolve973            FROM resolved_merchants974            WHERE outcome != 'unknown'975            GROUP BY outcome, issue_type976            ORDER BY count DESC977        """)978        rows = _fetch_all(conn, cursor)979 980        cursor.execute("""981            SELECT982                COUNT(*) as total_resolved,983                COUNT(CASE WHEN outcome = 'fixed' THEN 1 END) as fixed,984                COUNT(CASE WHEN outcome = 'false_positive' THEN 1 END) as false_positives,985                COUNT(CASE WHEN outcome = 'churned' THEN 1 END) as churned,986                COUNT(CASE WHEN outcome = 'closed' THEN 1 END) as closed,987                ROUND(AVG(daily_loss_at_resolution), 0) as avg_loss_saved988            FROM resolved_merchants989        """)990        summary = dict(cursor.fetchone())991 992        return {993            "summary": summary,994            "by_outcome_and_type": rows,995        }996    finally:997        conn.commit()998        release_connection(conn)999 1000def get_bank_scorecard(reference_date=None, days=7):1001    if reference_date:1002        cutoff = (reference_date - timedelta(days=days)).strftime("%Y-%m-%d")1003    else:1004        cutoff = (datetime.now() - timedelta(days=days)).strftime("%Y-%m-%d")1005    p = "%s"1006 1007    conn = get_connection()1008    try:1009        cursor = conn.cursor()1010        cursor.execute(f"""1011            SELECT m.acquiring_bank as bank, COUNT(*) as total_txns,1012                   SUM(CASE WHEN t.status = 'success' THEN 1 ELSE 0 END) as success_count,1013                   SUM(CASE WHEN t.status != 'success' THEN 1 ELSE 0 END) as failure_count,1014                   ROUND(AVG(t.bank_latency_ms), 0) as avg_bank_latency,1015                   ROUND(AVG(CASE WHEN t.status != 'success' THEN t.bank_latency_ms END), 0) as avg_fail_latency,1016                   ROUND(SUM(t.amount), 0) as total_amount,1017                   ROUND(SUM(CASE WHEN t.status = 'success' THEN t.amount ELSE 0 END), 0) as success_amount1018            FROM transactions t JOIN merchants m ON t.merchant_id = m.id WHERE t.day >= {p}1019            GROUP BY m.acquiring_bank ORDER BY failure_count DESC1020        """, (cutoff,))1021        banks = _fetch_all(conn, cursor)1022 1023        cursor.execute(f"""1024            SELECT m.acquiring_bank as bank, t.failure_reason, COUNT(*) as cnt, 1025                   ROUND(SUM(t.amount), 0) as failed_amount, ROUND(AVG(t.bank_latency_ms), 0) as avg_lat1026            FROM transactions t JOIN merchants m ON t.merchant_id = m.id1027            WHERE t.day >= {p} AND t.status != 'success' AND t.failure_reason IS NOT NULL1028            GROUP BY m.acquiring_bank, t.failure_reason ORDER BY m.acquiring_bank, cnt DESC1029        """, (cutoff,))1030        reason_rows = _fetch_all(conn, cursor)1031        top_reasons = {}1032        for r in reason_rows:1033            bank = r["bank"]1034            if bank not in top_reasons: top_reasons[bank] = {"reason": r["failure_reason"], "count": r["cnt"], "failed_amount": r["failed_amount"], "avg_latency": r["avg_lat"]}1035 1036        for b in banks:1037            b["success_rate"] = round(b["success_count"] / max(b["total_txns"], 1), 3)1038            b["failure_rate"] = round(b["failure_count"] / max(b["total_txns"], 1), 3)1039            b["top_failure_reason"] = top_reasons.get(b["bank"])1040        return {"period": f"last {days} days", "banks": banks}1041    finally:1042        conn.commit()1043        release_connection(conn)