import pandas as pd import numpy as np import requests import uvicorn import asyncio import re import uuid import json import io import zipfile from datetime import timedelta from fastapi import FastAPI, Body, HTTPException from fastapi.responses import HTMLResponse, JSONResponse, Response from huggingface_hub import HfApi, hf_hub_download from typing import Optional, List # ========================================== # 1. CONFIGURATION & SECRETS # ========================================== SB_URL = "https://txmakjdjvtwvdbwkeopd.supabase.co" SB_KEY = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJpc3MiOiJzdXBhYmFzZSIsInJlZiI6InR4bWFramRqdnR3dmRid2tlb3BkIiwicm9sZSI6InNlcnZpY2Vfcm9sZSIsImlhdCI6MTc4MTQ4ODU2OSwiZXhwIjoyMDk3MDY0NTY5fQ.OQTQ9M6SlCYGRZdQgmFQVHfwUhVKUQ0iBcLuitrHdok" HEADERS = { "apikey": SB_KEY, "Authorization": f"Bearer {SB_KEY}", "Content-Type": "application/json", "Prefer": "return=representation" } HF_TOKEN_PART_1 = "hf_EzaiwOHiejWFLN" HF_TOKEN_PART_2 = "otvQqTPUTksQQfYxrbJg" HF_TOKEN = HF_TOKEN_PART_1 + HF_TOKEN_PART_2 HF_REPO_ID = "topsecrettraders/DepthChain" hf_api = HfApi(token=HF_TOKEN) app = FastAPI() SESSION_CACHE = {} # ========================================== # 2. SUPABASE HELPER FUNCTIONS # ========================================== def sb_rpc(func_name: str, params: dict = None): if params is None: params = {} try: url = f"{SB_URL}/rest/v1/rpc/{func_name}" r = requests.post(url, headers=HEADERS, json=params) r.raise_for_status() return r.json() except Exception as e: err_msg = r.text if 'r' in locals() and hasattr(r, 'text') else str(e) print(f"DB Error RPC ({func_name}): {err_msg}") return [] # ========================================== # 3. API ROUTES - METADATA & UI # ========================================== @app.get("/") def home(): return HTMLResponse(HTML_TEMPLATE) @app.get("/api/dates") def get_dates(): data = sb_rpc("get_available_dates", {"p_limit": 50}) if isinstance(data, list): return JSONResponse([row['available_date'] for row in data if 'available_date' in row]) return JSONResponse([]) @app.get("/api/roots") def get_roots(date: str): data = sb_rpc("get_roots_for_date", {"p_date": date, "p_q": ""}) if isinstance(data, list): return JSONResponse([row['root_name'] for row in data if 'root_name' in row]) return JSONResponse([]) @app.get("/api/expiries") def get_expiries(date: str, root: str): data = sb_rpc("get_expiries_for_root", {"p_date": date, "p_root": root.upper()}) if isinstance(data, list): return JSONResponse([row['expiry_date'] for row in data if 'expiry_date' in row]) return JSONResponse([]) # ========================================== # 4. HUGGING FACE EXPLORER API # ========================================== @app.get("/api/hf_files") def list_hf_files(): try: files = hf_api.list_repo_files(repo_id=HF_REPO_ID, repo_type="dataset") sessions = {} for f in files: if f.startswith("Sessions/") and f.endswith("data.parquet"): parts = f.split("/") if len(parts) >= 4: session_name = parts[1] dataset_name = parts[2] if session_name not in sessions: sessions[session_name] = [] if dataset_name not in sessions[session_name]: sessions[session_name].append(dataset_name) # Sort sessions and datasets sorted_sessions = {k: sorted(v, reverse=True) for k, v in sorted(sessions.items(), reverse=True)} return JSONResponse({"status": "success", "sessions": sorted_sessions}) except Exception as e: return JSONResponse({"status": "error", "message": str(e)}) @app.get("/api/hf_metadata") def get_hf_metadata(session_name: str, dataset_name: str): try: meta_path = f"Sessions/{session_name}/{dataset_name}/metadata.json" local_path = hf_hub_download(repo_id=HF_REPO_ID, filename=meta_path, repo_type="dataset", token=HF_TOKEN) with open(local_path, "r") as f: metadata = json.load(f) return JSONResponse({"status": "success", "metadata": metadata}) except Exception as e: return JSONResponse({"status": "error", "message": "Metadata not found or error fetching."}) @app.get("/api/view_parquet") def view_parquet(session_name: str, dataset_name: str): try: file_path = f"Sessions/{session_name}/{dataset_name}/data.parquet" local_path = hf_hub_download(repo_id=HF_REPO_ID, filename=file_path, repo_type="dataset", token=HF_TOKEN) df = pd.read_parquet(local_path) # Limit rows to 1000 to prevent browser crash, it's just a viewer df_preview = df.head(1000) return JSONResponse({"status": "success", "columns": df_preview.columns.tolist(), "data": df_preview.to_dict(orient="records")}) except Exception as e: return JSONResponse({"status": "error", "message": str(e)}) @app.get("/api/download_session_zip") def download_session_zip(session_name: str): try: folder_path = f"Sessions/{session_name}" files = hf_api.list_repo_files(repo_id=HF_REPO_ID, repo_type="dataset") folder_files = [f for f in files if f.startswith(folder_path + "/")] if not folder_files: return JSONResponse({"status": "error", "message": "No files found for this session."}) zip_buffer = io.BytesIO() with zipfile.ZipFile(zip_buffer, "w", zipfile.ZIP_DEFLATED) as zip_file: for f_path in folder_files: local_path = hf_hub_download(repo_id=HF_REPO_ID, filename=f_path, repo_type="dataset", token=HF_TOKEN) # Keep the dataset folder structure inside the zip arcname = f_path.replace(f"{folder_path}/", "") zip_file.write(local_path, arcname=arcname) return Response( content=zip_buffer.getvalue(), media_type="application/zip", headers={"Content-Disposition": f'attachment; filename="{session_name}_Full.zip"'} ) except Exception as e: return JSONResponse({"status": "error", "message": str(e)}) @app.get("/api/download_session_parquets_zip") def download_session_parquets_zip(session_name: str): try: folder_path = f"Sessions/{session_name}" files = hf_api.list_repo_files(repo_id=HF_REPO_ID, repo_type="dataset") # Only grab the data.parquet files for this session parquet_files = [f for f in files if f.startswith(folder_path + "/") and f.endswith("data.parquet")] if not parquet_files: return JSONResponse({"status": "error", "message": "No Parquet files found for this session."}) zip_buffer = io.BytesIO() with zipfile.ZipFile(zip_buffer, "w", zipfile.ZIP_DEFLATED) as zip_file: for f_path in parquet_files: local_path = hf_hub_download(repo_id=HF_REPO_ID, filename=f_path, repo_type="dataset", token=HF_TOKEN) # Extract dataset name (e.g. 2026-07-10_NIFTY_2026-07-14) # Structure: Sessions / {session_name} / {dataset_name} / data.parquet parts = f_path.split("/") if len(parts) >= 4: dataset_name = parts[2] # Rename the file inside the zip to match requested format arcname = f"{dataset_name}_data.parquet" zip_file.write(local_path, arcname=arcname) return Response( content=zip_buffer.getvalue(), media_type="application/zip", headers={"Content-Disposition": f'attachment; filename="{session_name}_Parquets_Only.zip"'} ) except Exception as e: return JSONResponse({"status": "error", "message": str(e)}) # ========================================== # 5. STEP 1: FETCH & INSPECT # ========================================== @app.post("/api/inspect_data") async def inspect_data( date: str = Body(...), root: str = Body(...), expiry: str = Body(...), start_time: str = Body("09:15"), end_time: str = Body("15:30"), atm_range: int = Body(10), required_vars: List[str] = Body(["P", "CB", "PS", "CS", "PB"]) ): root = root.upper() t_start_dt = pd.to_datetime(f"{date} {start_time}:00") t_end_dt = pd.to_datetime(f"{date} {end_time}:00") inst_data = sb_rpc("get_instruments_for_expiry", {"p_date": date, "p_root": root, "p_expiry": expiry}) instruments = [r['instrument_name'] for r in inst_data if 'instrument_name' in r] if not instruments: return JSONResponse({"error": f"No data found in Supabase for {root} {expiry} on {date}"}) spot_keys =[k for k in instruments if k.endswith('-INDEX') or k.endswith('-EQ')] fut_keys =[k for k in instruments if 'FUT' in k] ref_instrument = spot_keys[0] if spot_keys else (fut_keys[0] if fut_keys else instruments[0]) strike_map = {} for inst in instruments: if inst.endswith('CE') or inst.endswith('PE'): clean_inst = inst.split(":")[-1] if clean_inst.startswith(root): match = re.search(r'(\d+(?:\.\d+)?)(?:CE|PE)$', clean_inst[len(root)+5:]) if match: strike_map[inst] = float(match.group(1)) continue matches = re.findall(r'(\d+(?:\.\d+)?)(?=CE$|PE$)', clean_inst) if matches: strike_map[inst] = float(matches[-1]) base_params = { "p_root": root, "p_expiry": expiry, "p_ref_instrument": ref_instrument, "p_atm_range": atm_range, "p_strike_map": strike_map } async def fetch_chunk(start_dt, end_dt, retries=3): params = base_params.copy() params["p_start_time"] = start_dt.strftime("%Y-%m-%d %H:%M:%S") params["p_end_time"] = end_dt.strftime("%Y-%m-%d %H:%M:%S") loop = asyncio.get_event_loop() for attempt in range(retries): try: return await loop.run_in_executor(None, sb_rpc, "get_advanced_chart_data", params) except Exception as e: if attempt == retries - 1: return[] await asyncio.sleep(1) tasks =[] curr = t_start_dt while curr <= t_end_dt: nxt = curr + timedelta(hours=1) - timedelta(seconds=1) if nxt > t_end_dt: nxt = t_end_dt tasks.append(fetch_chunk(curr, nxt)) curr += timedelta(hours=1) chunk_results = await asyncio.gather(*tasks) all_records =[] for res in chunk_results: if isinstance(res, list): all_records.extend(res) if not all_records: return JSONResponse({"error": "No valid records fetched. Data might be empty for this date/time."}) df = pd.DataFrame(all_records) df.rename(columns={"tick_ts": "timestamp", "price": "P", "cb": "CB", "cs": "CS", "pb": "PB", "ps": "PS", "b_qty": "B", "s_qty": "S", "vol": "V", "oi": "OI"}, inplace=True) df['timestamp'] = pd.to_datetime(df['timestamp']) df.drop_duplicates(subset=['timestamp'], inplace=True) df.set_index('timestamp', inplace=True) df.sort_index(inplace=True) df = df.loc[(df.index >= t_start_dt) & (df.index <= t_end_dt)] if df.empty: return JSONResponse({"error": "Data exists, but completely outside your specified Start/End Time."}) available_vars =[v for v in required_vars if v in df.columns] df = df[available_vars].copy() actual_min = df.index.min() actual_max = df.index.max() full_idx = pd.date_range(start=actual_min, end=actual_max, freq='10s') raw_ticks = len(df) expected_ticks = len(full_idx) missing_ts = full_idx.difference(df.index).to_pydatetime() missing_blocks =[] if len(missing_ts) > 0: start_block = missing_ts[0] prev = missing_ts[0] for ts in missing_ts[1:]: if (ts - prev).total_seconds() > 10.5: missing_blocks.append({ "start": start_block.strftime('%Y-%m-%d %H:%M:%S'), "end": prev.strftime('%Y-%m-%d %H:%M:%S'), "count": int((prev - start_block).total_seconds()/10) + 1 }) start_block = ts prev = ts missing_blocks.append({ "start": start_block.strftime('%Y-%m-%d %H:%M:%S'), "end": prev.strftime('%Y-%m-%d %H:%M:%S'), "count": int((prev - start_block).total_seconds()/10) + 1 }) df_reindexed = df.reindex(full_idx) null_details_existing = {} for col in available_vars: na_times = df_reindexed[df_reindexed[col].isna() | (df_reindexed[col] == 0)].index.strftime('%H:%M:%S').tolist() null_details_existing[col] = {"count": len(na_times), "times": na_times[:50]} session_id = str(uuid.uuid4()) SESSION_CACHE[session_id] = df.copy() return JSONResponse({ "status": "success", "session_id": session_id, "stats": { "expected_ticks": expected_ticks, "found_ticks": raw_ticks, "missing_time_ticks": len(missing_ts), "missing_blocks": missing_blocks, "null_breakdown": null_details_existing, "variables": available_vars, "original_start": actual_min.strftime('%Y-%m-%d %H:%M:%S'), "original_end": actual_max.strftime('%Y-%m-%d %H:%M:%S') }, "meta": { "date": date, "root": root, "expiry": expiry, "start": start_time, "end": end_time, "atm_range": atm_range } }) # ========================================== # 6. STEP 2: CLEAN & ARCHIVE # ========================================== @app.post("/api/clean_and_upload") async def clean_and_upload( session_id: str = Body(...), session_name: str = Body(...), fill_method: str = Body("ffill"), trim_start: str = Body(None), trim_end: str = Body(None), metadata: dict = Body(...) ): if session_id not in SESSION_CACHE: return JSONResponse({"error": "Session expired or invalid. Please Fetch & Inspect again."}) # Validate session name clean_session_name = re.sub(r'[^a-zA-Z0-9_\-]', '_', session_name) if not clean_session_name: clean_session_name = "Default_Session" df = SESSION_CACHE[session_id].copy() t_start = pd.to_datetime(trim_start) if trim_start else df.index.min() t_end = pd.to_datetime(trim_end) if trim_end else df.index.max() df = df.loc[(df.index >= t_start) & (df.index <= t_end)] full_idx = pd.date_range(start=t_start, end=t_end, freq='10s') df = df.reindex(full_idx) df.replace(0, np.nan, inplace=True) if fill_method == "ffill": df_clean = df.ffill().fillna(0) elif fill_method == "interpolate": df_clean = df.interpolate(method='linear').fillna(0) elif fill_method == "zero": df_clean = df.fillna(0) else: df_clean = df.ffill().fillna(0) df_clean.reset_index(inplace=True) df_clean.rename(columns={"index": "timestamp"}, inplace=True) df_clean['timestamp'] = df_clean['timestamp'].dt.strftime('%Y-%m-%d %H:%M:%S') metadata["session_name"] = clean_session_name metadata["imputation_applied"] = fill_method metadata["trim_start_applied"] = trim_start metadata["trim_end_applied"] = trim_end metadata["total_final_rows"] = len(df_clean) date = metadata["date"] root = metadata["root"] expiry = metadata["expiry"] target_ds_name = f"{date}_{root}_{expiry}" # Handle duplicates inside the specific session folder try: files = hf_api.list_repo_files(repo_id=HF_REPO_ID, repo_type="dataset") session_files = [f for f in files if f.startswith(f"Sessions/{clean_session_name}/")] final_ds_name = target_ds_name version = 1 while any(f.startswith(f"Sessions/{clean_session_name}/{final_ds_name}/") for f in session_files): version += 1 final_ds_name = f"{target_ds_name}_v{version}" hf_folder = f"Sessions/{clean_session_name}/{final_ds_name}" except Exception as e: hf_folder = f"Sessions/{clean_session_name}/{target_ds_name}" try: parquet_buffer = io.BytesIO() df_clean.to_parquet(parquet_buffer, index=False) parquet_buffer.seek(0) json_buffer = io.BytesIO(json.dumps(metadata, indent=4).encode('utf-8')) json_buffer.seek(0) hf_api.upload_file(path_or_fileobj=parquet_buffer, path_in_repo=f"{hf_folder}/data.parquet", repo_id=HF_REPO_ID, repo_type="dataset") hf_api.upload_file(path_or_fileobj=json_buffer, path_in_repo=f"{hf_folder}/metadata.json", repo_id=HF_REPO_ID, repo_type="dataset") del SESSION_CACHE[session_id] return JSONResponse({"status": "success", "message": f"Dataset safely archived to {hf_folder}!"}) except Exception as e: return JSONResponse({"error": f"Upload Failed: {str(e)}"}) # ========================================== # 7. FRONTEND TEMPLATE (LIGHT MODE) # ========================================== HTML_TEMPLATE = """
All datasets extracted will be saved inside this session folder.