import pandas as pd import requests import uvicorn import asyncio import re from datetime import timedelta from fastapi import FastAPI, Body from fastapi.responses import HTMLResponse, JSONResponse from typing import Optional, List # ========================================== # 1. CONFIGURATION # ========================================== SB_URL = "https://ivyyfhymhdheykpntozd.supabase.co" SB_KEY = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJpc3MiOiJzdXBhYmFzZSIsInJlZiI6Iml2eXlmaHltaGRoZXlrcG50b3pkIiwicm9sZSI6InNlcnZpY2Vfcm9sZSIsImlhdCI6MTc4Mzg5NzAwOCwiZXhwIjoyMDk5NDczMDA4fQ.cHvePFN3NML8Zv12D2vV91MigJGRerZvD_Cevbu1AFk" HEADERS = { "apikey": SB_KEY, "Authorization": f"Bearer {SB_KEY}", "Content-Type": "application/json", "Prefer": "return=representation" } app = FastAPI() # ========================================== # 2. HELPER FUNCTIONS # ========================================== def sb_get(endpoint: str, params: dict = None): try: url = f"{SB_URL}/rest/v1{endpoint}" r = requests.get(url, headers=HEADERS, params=params) r.raise_for_status() return r.json() except Exception as e: print(f"DB Error GET ({endpoint}): {e}") return [] 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: print(f"DB Error RPC ({func_name}): {e}") return [] # ========================================== # 3. API ROUTES - METADATA # ========================================== @app.get("/") def home(): return HTMLResponse(HTML_TEMPLATE) @app.get("/api/dates") def get_dates(): data = sb_rpc("get_available_dates", {"p_limit": 30}) dates = [row['available_date'] for row in data] return JSONResponse(dates) @app.get("/api/roots") def get_roots(date: str, q: Optional[str] = None): params = {"p_date": date} if q: params["p_q"] = q.upper() data = sb_rpc("get_roots_for_date", params) roots = [row['root_name'] for row in data] return JSONResponse(roots) @app.get("/api/auto_config") def get_auto_config(date: str, root: str): root = root.upper() exps_data = sb_rpc("get_expiries_for_root", {"p_date": date, "p_root": root}) all_expiries = [row['expiry_date'] for row in exps_data if row['expiry_date'] != 'MARKET'] valid_expiries = [e for e in all_expiries if e >= date] if not valid_expiries: valid_expiries = all_expiries response = {"expiries": all_expiries, "current": None, "next": None} def find_fut_rpc(exp): inst_data = sb_rpc("get_instruments_for_expiry", {"p_date": date, "p_root": root, "p_expiry": exp}) instruments = [r['instrument_name'] for r in inst_data] for i in instruments: if 'FUT' in i: return i return instruments[0] if instruments else None if len(valid_expiries) > 0: exp_c = valid_expiries[0] response["current"] = {"expiry": exp_c, "instrument": find_fut_rpc(exp_c)} if len(valid_expiries) > 1: exp_n = valid_expiries[1] response["next"] = {"expiry": exp_n, "instrument": find_fut_rpc(exp_n)} return JSONResponse(response) @app.get("/api/expiries") def get_expiries(date: str, root: str): root = root.upper() data = sb_rpc("get_expiries_for_root", {"p_date": date, "p_root": root}) exps = [row['expiry_date'] for row in data] return JSONResponse(exps) @app.get("/api/instruments") def get_instruments(date: str, root: str, expiry: str): root = root.upper() data = sb_rpc("get_instruments_for_expiry", {"p_date": date, "p_root": root, "p_expiry": expiry}) instruments = [row['instrument_name'] for row in data] spot_fut = [] options = [] for i in instruments: if i.endswith('CE') or i.endswith('PE'): options.append(i) else: spot_fut.append(i) return JSONResponse({"spot_fut": sorted(spot_fut), "options": sorted(options)}) # ========================================== # 4. API ROUTES - DATA FETCHING # ========================================== @app.post("/api/fetch_series") async def fetch_series( date: str = Body(...), root: str = Body(...), expiry: str = Body(...), instrument: str = Body(""), timeframe: str = Body(...), start_time: str = Body(...), end_time: str = Body(...), mode: str = Body("normal"), atm_range: int = Body(5), required_vars: List[str] = Body(default=None), ref_times: List[str] = Body(default=None) ): root = root.upper() t_start_dt = pd.to_datetime(f"{date} {start_time}:00") t_end_dt = pd.to_datetime(f"{date} {end_time}:59") if required_vars is None: required_vars = ["P", "B", "S", "CB", "CS", "PB", "PS", "V", "OI"] if ref_times is None: ref_times = [] required_vars = list(set(required_vars) | {"P"}) if mode == "advanced": 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 not instruments: return JSONResponse({"error": f"No data found for expiry {expiry}", "labels": [], "P": []}) ref_instrument = instrument if not ref_instrument: 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] if spot_keys: ref_instrument = spot_keys[0] elif fut_keys: ref_instrument = fut_keys[0] else: others = [k for k in instruments if not k.endswith('CE') and not k.endswith('PE')] ref_instrument = others[0] if others else instruments[0] # Robust Python Logic for Strike Pricing mapping strike_map = {} for inst in instruments: if inst.endswith('CE') or inst.endswith('PE'): clean_inst = inst.split(":")[-1] if clean_inst.startswith(root): remainder = clean_inst[len(root):] match = re.search(r'(\d+(?:\.\d+)?)(?:CE|PE)$', remainder[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]) rpc_name = "get_advanced_chart_data" base_params = { "p_root": root, "p_expiry": expiry, "p_ref_instrument": ref_instrument, "p_atm_range": atm_range, "p_strike_map": strike_map } else: if not instrument: return JSONResponse({"error": "No Instrument Selected.", "labels": [], "P": []}) rpc_name = "get_normal_chart_data" base_params = { "p_root": root, "p_expiry": expiry, "p_instrument": instrument } # ========================================================= # FETCH MAIN DATA CHUNKS # ========================================================= async def fetch_chunk(start_dt, end_dt): 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() return await loop.run_in_executor(None, sb_rpc, rpc_name, params) 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 data found in range.", "labels": [], "P": []}) df = pd.DataFrame(all_records) df.rename(columns={ "tick_ts": "ts", "b_qty": "B", "s_qty": "S", "vol": "V", "oi": "OI", "cb": "CB", "cs": "CS", "pb": "PB", "ps": "PS", "price": "P" }, inplace=True) df['ts'] = pd.to_datetime(df['ts']) df.drop_duplicates(subset=['ts'], inplace=True) df.set_index('ts', inplace=True) df.sort_index(inplace=True) tf_map = { "1min": "1min", "3min": "3min", "5min": "5min", "15min": "15min", "30min": "30min", "1hour": "1h" } panda_tf = tf_map.get(timeframe, "1min") agg_dict = { "P": "last", "V": "sum", "OI": "last", "B": "last", "S": "last", "CB": "last", "CS": "last", "PB": "last", "PS": "last", } if mode == "advanced": agg_dict["strike_breakdown"] = "last" existing_agg = {k: v for k, v in agg_dict.items() if k in df.columns} resampled = df.resample(panda_tf).agg(existing_agg).ffill().fillna(0) # ========================================================= # FETCH REFERENCE TIMES # ========================================================= ref_results = {} if ref_times: async def fetch_and_agg_ref(rt): rt_start = pd.to_datetime(f"{date} {rt}:00") rt_end = pd.to_datetime(f"{date} {rt}:59") res = await fetch_chunk(rt_start, rt_end) if not res: return rt, None dfr = pd.DataFrame(res) dfr.rename(columns={ "tick_ts": "ts", "b_qty": "B", "s_qty": "S", "vol": "V", "oi": "OI", "cb": "CB", "cs": "CS", "pb": "PB", "ps": "PS", "price": "P" }, inplace=True) dfr['ts'] = pd.to_datetime(dfr['ts']) dfr.drop_duplicates(subset=['ts'], inplace=True) dfr.set_index('ts', inplace=True) dfr.sort_index(inplace=True) agg_d = {k: v for k, v in agg_dict.items() if k in dfr.columns and k != "strike_breakdown"} r_resampled = dfr.resample('1min').agg(agg_d).ffill().fillna(0) if r_resampled.empty: return rt, None return rt, r_resampled.iloc[0].to_dict() tasks_ref = [fetch_and_agg_ref(rt) for rt in set(ref_times)] ref_resolved = await asyncio.gather(*tasks_ref) for rt, row in ref_resolved: if row: ref_results[rt] = row response_data = { "labels": resampled.index.strftime('%H:%M').tolist(), "ref_values": ref_results } for var_key in ["P", "B", "S", "V", "OI", "CB", "CS", "PB", "PS"]: if var_key in required_vars and var_key in resampled.columns: response_data[var_key] = [round(float(v), 2) for v in resampled[var_key].tolist()] if mode == "advanced" and "strike_breakdown" in resampled.columns: response_data["strike_breakdown"] = [ x if isinstance(x, list) else [] for x in resampled["strike_breakdown"].tolist() ] return JSONResponse(response_data) # ========================================== # 5. FRONTEND TEMPLATE # ========================================== HTML_TEMPLATE = """