Prathmesh / app.py
parthislive's picture
Update app.py
83a181a verified
Raw History Blame Contribute Delete
8.07 kB
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=<num>": "ICMR by phone",
"/api/aadhar?aadhar=<num>": "ICMR by aadhar",
"/api/name?name=<name>": "ICMR by name",
"/api/user?user_id=<id>": "Search HiTeck + Telegram by user_id",
"/api/lookup?phone=<num>": "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)