mrik8899/fintech-demo
0
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)