from fastapi import FastAPI, Query, HTTPException import duckdb import time import os import threading app = FastAPI(title="ICMR + HiTeck + Telegram API") # ---------- ICMR SORTED INDEX FILES ---------- HF_BASE_ICMR = "https://huggingface.co/datasets/parthislive/Ha/resolve/main" PHONE_FILES = ", ".join([f"'{HF_BASE_ICMR}/idx_phone.{i}.parquet'" for i in range(7)]) AADHAR_FILES = ", ".join([f"'{HF_BASE_ICMR}/idx_aadhar.{i}.parquet'" for i in range(7)]) DISPLAY_COLS = "name, fathersName, phoneNumber, aadharNumber, otherNumber, address, district, pincode, state, town, source" # ---------- HITECKGROUP DATASET ---------- HF_BASE_HITECK = "https://huggingface.co/datasets/parthislive/HiTeckGroup/resolve/main" HITECK_PARQUET = f"'{HF_BASE_HITECK}/data/train-00000-of-00001.parquet'" HITECK_COLS = "user_id, phone_number, country, country_code" # ---------- TELEGRAM DATASET ---------- HF_BASE_TG = "https://huggingface.co/datasets/parthislive/TELEGRAM/resolve/main" TG_PARQUET = f"'{HF_BASE_TG}/simple_all/merged_all.parquet'" TG_COLS = "user_id, phone_number, country, country_code" _con = None _lock = threading.Lock() def get_con(): global _con if _con is None: with _lock: if _con is None: con = duckdb.connect() con.execute("SET home_directory='/tmp'") con.execute("SET extension_directory='/tmp/duckdb_extensions'") con.execute("INSTALL httpfs; LOAD httpfs;") con.execute("INSTALL parquet; LOAD parquet;") con.execute("SET enable_object_cache=true") con.execute("SET threads=8") con.execute("SET http_keep_alive=true") con.execute("SET http_timeout=120000") _con = con return _con # ============================================================ # ROOT & HEALTH # ============================================================ @app.get("/") def root(): return { "app": "ICMR + HiTeck + Telegram API", "endpoints": { "/api/phone?phone=": "ICMR by phone", "/api/aadhar?aadhar=": "ICMR by aadhar", "/api/name?name=": "ICMR by name", "/api/user?user_id=": "Search HiTeck + Telegram by user_id", "/api/lookup?phone=": "Search HiTeck + Telegram by phone", } } @app.get("/health") def health(): return {"status": "ok"} # ============================================================ # ICMR ENDPOINTS # ============================================================ @app.get("/api/phone") def api_phone(phone: str = Query(...)): return _search_phone(phone.strip()) @app.get("/api/aadhar") def api_aadhar(aadhar: str = Query(...)): return _search_aadhar(aadhar.strip()) @app.get("/api/name") def api_name(name: str = Query(...)): return _search_generic("name", name.strip(), "contains", 20) def _search_phone(phone: str): if not phone: raise HTTPException(400, "Empty phone") try: start = time.time() con = get_con() sql = f""" SELECT {DISPLAY_COLS} FROM read_parquet([{PHONE_FILES}]) WHERE phoneNumber = ? LIMIT 10 """ cur = con.execute(sql, [phone]) cols = [d[0] for d in cur.description] rows = cur.fetchall() results = [dict(zip(cols, r)) for r in rows] elapsed = round(time.time() - start, 2) if not results: raise HTTPException(404, f"No data for phone={phone}") return {"status": "success", "field": "phoneNumber", "query": phone, "count": len(results), "time": f"{elapsed}s", "data": results} except HTTPException: raise except Exception as e: raise HTTPException(500, str(e)) def _search_aadhar(aadhar: str): if not aadhar: raise HTTPException(400, "Empty aadhar") try: start = time.time() con = get_con() sql = f""" SELECT {DISPLAY_COLS} FROM read_parquet([{AADHAR_FILES}]) WHERE aadharNumber = ? LIMIT 10 """ cur = con.execute(sql, [aadhar]) cols = [d[0] for d in cur.description] rows = cur.fetchall() results = [dict(zip(cols, r)) for r in rows] elapsed = round(time.time() - start, 2) if not results: raise HTTPException(404, f"No data for aadhar={aadhar}") return {"status": "success", "field": "aadharNumber", "query": aadhar, "count": len(results), "time": f"{elapsed}s", "data": results} except HTTPException: raise except Exception as e: raise HTTPException(500, str(e)) def _search_generic(field, value, mode, limit): try: start = time.time() con = get_con() all_files = f"{PHONE_FILES}, {AADHAR_FILES}" if mode == "exact": sql = f"SELECT {DISPLAY_COLS} FROM read_parquet([{all_files}]) WHERE {field} = ? LIMIT {limit}" params = [value] else: sql = f"SELECT {DISPLAY_COLS} FROM read_parquet([{all_files}]) WHERE {field} ILIKE ? LIMIT {limit}" params = [f"%{value}%"] cur = con.execute(sql, params) cols = [d[0] for d in cur.description] rows = cur.fetchall() results = [dict(zip(cols, r)) for r in rows] elapsed = round(time.time() - start, 2) if not results: raise HTTPException(404, f"No data for {field}={value}") return {"status": "success", "field": field, "query": value, "count": len(results), "time": f"{elapsed}s", "data": results} except HTTPException: raise except Exception as e: raise HTTPException(500, str(e)) # ============================================================ # COMBINED: HiTeck + Telegram (single endpoint each for user_id and phone) # ============================================================ def _combined_search(field: str, value: str): if not value: raise HTTPException(400, f"Empty {field}") start = time.time() con = get_con() hiteck_data = [] tg_data = [] errors = {} try: cur = con.execute( f"SELECT {HITECK_COLS} FROM read_parquet({HITECK_PARQUET}) WHERE {field} = ? LIMIT 10", [value], ) cols = [d[0] for d in cur.description] hiteck_data = [dict(zip(cols, r)) for r in cur.fetchall()] except Exception as e: errors["hiteck"] = str(e) try: cur = con.execute( f"SELECT {TG_COLS} FROM read_parquet({TG_PARQUET}) WHERE {field} = ? LIMIT 10", [value], ) cols = [d[0] for d in cur.description] tg_data = [dict(zip(cols, r)) for r in cur.fetchall()] except Exception as e: errors["telegram"] = str(e) elapsed = round(time.time() - start, 2) total = len(hiteck_data) + len(tg_data) if total == 0 and errors: raise HTTPException(500, errors) if total == 0: raise HTTPException(404, f"No data for {field}={value} in any dataset") return { "status": "success", "field": field, "query": value, "total_count": total, "time": f"{elapsed}s", "hiteck": {"count": len(hiteck_data), "data": hiteck_data}, "telegram": {"count": len(tg_data), "data": tg_data}, "errors": errors if errors else None, } @app.get("/api/user") def api_user(user_id: str = Query(...)): """Search by user_id across HiTeck + Telegram.""" return _combined_search("user_id", user_id.strip()) @app.get("/api/lookup") def api_lookup(phone: str = Query(...)): """Search by phone across HiTeck + Telegram.""" return _combined_search("phone_number", phone.strip()) # ============================================================ # RUN # ============================================================ if __name__ == "__main__": import uvicorn port = int(os.environ.get("PORT", 7860)) uvicorn.run(app, host="0.0.0.0", port=port)