import asyncio import aiohttp import pandas as pd import gradio as gr import os from datetime import datetime from pymongo import MongoClient import urllib wallets = [ "79x5TFJ1cNrPXaRVGeybfyqvzYfrBZKE7SasxAhacP6o", "88HqLVZTVaXdWLTykPvSNFzo7svunmqvJfbZrxURyDi7", "EMroWV2SJHWGTMnTvo9azmm29gS2g8NN4KkKtQMEJcZY", "34jPBC2QFYE3jefay7czTZuSnBv1sghGr8nkBuS6tvpw", "6VsmE36FQmy8dyoqpvmZVn9iRiNnpi7zG9G2sQy6z2Cs", "Ag7s3n6jrc4fxU78UwncsvNwuCbhG9EWmagrwcRJLfk6", "8fbyZjwfMgYmtaGcvMFmLKVciNgtWC43cfWp3r77AKuK", "FDQs4iWM9i5wVHkkHwZV3GAAZj2kHECM2jKtUi7SPweb", "FzbbDiX6X5kYJksjr9WcJG7Un7xr1sz3TjKF4KQKVVPX", "8y9Zd32B3Dpmd62Jnjv8zqfxomvNsMmJ2PHUoLLrSgZR", "Ce5MTdwh5RwtkqPydSUxXTmiVquytMEvZEzbB4Jx6KTa", "5bzj6Dc4e45UHrbBoRAKBy9TT24KzY5tGGvVAqZQ13M7", "8sf7ALGiFxjDPnHJWmbvicf68kqjfXtCrVAbMiuumhBh", "2uuyc8wYpwsSP3wdYiq2P3cZWKfm29QwJK5wrbQ4ZCzn", "6FmvFm112eKmVCzayD7LM1taYnaHr3qTscXhA3dWGRUw", "G83K6FiRsymtxTFRtxKSwgW5JMuHiU9eyNsPVvy7Gr77", "Et6HKKxRjQ17Mjaak1zL3rEtEchq2vwxEKWg46ucm9gW", "4tin3hjY3iFpLYzRU23gAojiZYugfUiZHk5mjVCdXWpk", "3TSkmRjoBvUqTDNc4z6FiPstmXJLavNsMTaprdJgRfUk", "C7ndjsaetCPES7A8PfzNwTeuSMXBgfKFg9KQLWiMHpLs", "6BZyeqvnRLhcFA8dNRQgi7xsccqEmJg8Dni1rn8u89xt", "9MVvCipSPaXXthTRH8zRETTF2e2gxSSgtKstfMyac96T", "DAkUEsfx7QuxofA7ejRf6JGkGMk3gwsQbauhaADZFvru", "8HyXT9MCAsmga4iGJeEBXURoDHstjdAnhdLK46mtth9E", "HQ7JpuRVdABJH3TjB4LhQVNQNzpQdsvzbEJHoiHEFUFQ", "E1bH8zEio6uvfFAagpr6eg67kFrSoXnaQbyQhyQJRtbg", "4987e4gBT9ZKNEpsRnt2kax9qT1V5VUj9MPoAtdQuZ4X", "8LAA9CCMuD5A9gHmGvcpfAMQLQt4VKkaDHJaD3oUM6Dq", "3sP4DM1eQCWdNeVkfxEqXkzVoAhEonTiNzSV5Nc8ctuY", "A9Pb6abhVBnmQeE5Hn3p1KdifegeiDNhgMbLoF7fGpgS", "3Rq3vDedATHBmB2t811VsGfqvD5jmNSHjCWf9pxpL9pe", "5uD6sBkFX6byxHeL9YSige1p2bZnHKnNWaCYF3d814wh", "C3t6E3rUyBEoJBRbYFTV4jrmaK2mkZPFjbVYymv3aLrK", "NFT85KM4R3bB83kGfFiaCxiptcVNSkNAvNmzXKVv8FR", "CBmQQDRiREgpUFUpwCu7bhpj6w66vPDKnjL4BxLxAi22", "5DXH2RV6CYXziBD2k9HN9kqukhF1A195xvvbtSRWUwMJ" ] API_TOKEN = os.getenv('SOLSCAN_KEY') API_TOKEN = 'eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJjcmVhdGVkQXQiOjE3NDM2OTQ5OTY4OTcsImVtYWlsIjoiamFtYWFsLm11dGFsaXBoQHRyaXJlbWV0cmFkaW5nLmNvbSIsImFjdGlvbiI6InRva2VuLWFwaSIsImFwaVZlcnNpb24iOiJ2MiIsImlhdCI6MTc0MzY5NDk5Nn0.7qakaGTx9vPod2zjbaulgPV3kXILQABF0XokxE-ocPA' headers = { "token": API_TOKEN } async def fetch_wallet_data(session, wallet): url = f"https://pro-api.solscan.io/v2.0/account/defi/activities?address={wallet}&activity_type[]=ACTIVITY_TOKEN_SWAP&activity_type[]=ACTIVITY_AGG_TOKEN_SWAP&page=1&page_size=100&sort_by=block_time&sort_order=desc" try: async with session.get(url) as response: data = await response.json() if 'data' in data: df = pd.json_normalize(data['data']) df['wallet'] = wallet return df return None except Exception as e: print(f"Error fetching data for wallet {wallet}: {e}") return None async def process(): async with aiohttp.ClientSession(headers=headers) as session: tasks = [fetch_wallet_data(session, wallet) for wallet in wallets] results = await asyncio.gather(*tasks) # Filter out None results dataframes = [df for df in results if df is not None] if dataframes: full_df = pd.concat(dataframes, ignore_index=True) # Drop nested column if it exists full_df = full_df.drop(columns=['routers.child_routers'], errors='ignore') return full_df else: print("No data fetched.") return pd.DataFrame() def get_token_amount(row): if row['token in'] == '6D6ccmg71x56V5Je1Mh82MFPYL38gaZqNc2LG1XMbonk': return row['token in amount'] elif row['token out'] == '6D6ccmg71x56V5Je1Mh82MFPYL38gaZqNc2LG1XMbonk': return row['token out amount'] else: return None # or np.nan def classify_trade_direction(row): if row['token in'] == '6D6ccmg71x56V5Je1Mh82MFPYL38gaZqNc2LG1XMbonk': return 'sell' elif row['token out'] == '6D6ccmg71x56V5Je1Mh82MFPYL38gaZqNc2LG1XMbonk': return 'buy' else: return None # or 'other', if you prefer def get_dollar_delta(row): if row['type'] == 'sell': return row['value'] elif row['type'] == 'buy': return -row['value'] else: return None # or 'other', if you prefer def get_token_delta(row): if row['type'] == 'sell': return -row['base_token_amount'] elif row['type'] == 'buy': return row['base_token_amount'] else: return None # or 'other', if you prefer #LIMIT ORDERS =================================================================================== SCRAPE_DO_TOKEN = os.getenv('SCRAPE_DO_TOKEN') client = MongoClient(os.getenv('MONGO_DB_URI')) db_1 = client["billy_balances"] token_collection = db_1["turnstile-tokens"] turnstile_token = token_collection.find_one(sort=[("timestamp", -1)])['x-turnstile-token'] common_headers = { "accept": "application/json", "accept-encoding": "identity", "accept-language": "en-GB,en-US;q=0.9,en;q=0.8", "authorization": "Bearer CGtF4EdvDbBpwUXmZSKW3HsYkajy7e", "content-type": "application/json", "origin": "https://portfolio.jup.ag", "referer": "https://portfolio.jup.ag/", "user-agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/136.0.0.0 Safari/537.36", "x-turnstile-token": turnstile_token # keep this updated } async def fetch_portfolio(session, wallet_address, results): base_url = f"https://portfolio-api-jup.sonar.watch/v1/portfolio/fetch?address={wallet_address}&addressSystem=solana" encoded_url = urllib.parse.quote(base_url) proxy_url = f"http://api.scrape.do/?token={SCRAPE_DO_TOKEN}&url={encoded_url}&forwardHeaders=True&super=True" try: async with session.get(proxy_url, headers=common_headers) as resp: text = await resp.json() if resp.status == 200: print(f"✅ Wallet {wallet_address[:6]}...: success") results.append({"wallet": wallet_address,"data": text}) else: print(f"❌ Wallet {wallet_address[:6]}...: HTTP {resp.status}") results.append({"wallet": wallet_address, "data": None}) except Exception as e: print(f"⚠️ Error fetching {wallet_address}: {e}") results.append({"wallet": wallet_address, "status": "error", "error": str(e)}) BATCH_SIZE = 20 async def get_data(wallet_addresses): results = [] async with aiohttp.ClientSession() as session: # Split addresses into batches for i in range(0, len(wallet_addresses), BATCH_SIZE): batch = wallet_addresses[i:i + BATCH_SIZE] tasks = [fetch_portfolio(session, address, results) for address in batch] await asyncio.gather(*tasks) await asyncio.sleep(0.5) # optional: rate limit delay return results all_elements = [] async def process_limits(): portfolios = await get_data(wallets) df = pd.DataFrame(portfolios) print(df) for index, row in df.iterrows(): data = row['data'] original_row = row['wallet'] # Token info mapping token_info = data.get('tokenInfo', {}).get('solana', {}) # Navigate to assets list elements = data.get('elements', []) for element in elements: if element.get("platformId") == "jupiter-exchange": element['wallet'] = original_row all_elements.append(element) '''input_token = element.get("data", {}).get("assets", {}).get("input", {}) if input_token: token_data = input_token.get("data", {}) address = token_data.get("address") value = input_token.get("value") amount = token_data.get("amount") symbol = token_info.get(address, {}).get("symbol", address) jup_rows.append({ "original_row": original_row, "address": address, "symbol": symbol, "amount": amount, "value": value })''' async def process_limitss(): await process_limits() order_df = pd.json_normalize(all_elements) print(order_df) def classify_type(row): if row['data.inputAddress'] == '6D6ccmg71x56V5Je1Mh82MFPYL38gaZqNc2LG1XMbonk': return 'sell' elif row['data.outputAddress'] == '6D6ccmg71x56V5Je1Mh82MFPYL38gaZqNc2LG1XMbonk': return 'buy' else: return None # or 'other', if you prefer order_df['type'] = order_df.apply(classify_type, axis=1) def get_token_amounts(row): if row['type'] == 'sell': return float(row['data.assets.input.data.amount']) elif row['type'] == 'buy': return float(row['data.expectedOutputAmount']) else: return None # or 'other', if you prefer order_df['base_token_amount'] = order_df.apply(get_token_amounts, axis=1) order_df['outputValue'] = order_df['data.outputPrice'] * order_df['data.expectedOutputAmount'] def get_price(row): if row['type'] == 'buy': return row['value']/row['base_token_amount'] elif row['type'] == 'sell': return row['outputValue']/row['base_token_amount'] else: return None def get_value(row): if row['type'] == 'buy': return row['value'] elif row['type'] == 'sell': return row['outputValue'] order_df['price'] = order_df.apply(get_price,axis=1) order_df['value'] = order_df.apply(get_value,axis=1) order_df = order_df[['wallet','label','type','value','price','data.inputAddress','data.outputAddress','data.filledPercentage']] order_df = order_df.rename(columns={'label':'order type','data.inputAddress':'token in','data.outputAddress':'token out','data.filledPercentage':'filled percentage'}) order_df = order_df.rename(columns={'type':'side','order type':'type','value':'value($)'}) order_df = order_df.sort_values('price',ascending=False).reset_index(drop=True) order_df['wallet number'] = order_df['wallet'].apply(lambda x : wallets.index(x) + 1) order_df = order_df[[order_df.columns[-1]] + list(order_df.columns[:-1])] return order_df # Setup sync wrapper async def main(): # Run async main and prepare data # Run the async code final_df = await process() order_df = await process_limitss() cols_to_convert = ['routers.amount1', 'routers.token1_decimals'] final_df[cols_to_convert] = final_df[cols_to_convert].apply(pd.to_numeric, errors='coerce') cols_to_convert = ['routers.amount2', 'routers.token2_decimals'] final_df[cols_to_convert] = final_df[cols_to_convert].apply(pd.to_numeric, errors='coerce') final_df['routers.amount1'] = final_df['routers.amount1'] / (10 ** final_df['routers.token1_decimals']) final_df['routers.amount2'] = final_df['routers.amount2'] / (10 ** final_df['routers.token2_decimals']) final_df.drop(columns=['routers.token1_decimals','routers.token2_decimals','platform','sources','activity_type','from_address'],inplace=True) final_df.drop(columns=['block_id','block_time'],inplace=True) final_df.rename(columns={'routers.token1':'token in','routers.token2':'token out','routers.amount1':'token in amount','routers.amount2':'token out amount'},inplace=True) final_df['wallet number'] = final_df['wallet'].apply(lambda x : wallets.index(x) + 1) final_df['type'] = final_df.apply(classify_trade_direction, axis=1) final_df['base_token_amount'] = final_df.apply(get_token_amount, axis=1) final_df['price'] = final_df['value']/final_df['base_token_amount'] final_df = final_df[['wallet number','wallet','time','type','value','price','trans_id','base_token_amount']] final_df = final_df.rename(columns={"trans_id":"tx_hash"}) final_df['dollar delta'] = final_df.apply(get_dollar_delta, axis=1) final_df['token delta'] = final_df.apply(get_token_delta, axis=1) final_df = final_df[['wallet number','wallet','time','type','price','dollar delta','token delta','tx_hash']] final_df.sort_values('time',ascending=False,inplace=True) final_df['time'] = final_df['time'].apply( lambda x: datetime.fromisoformat(x.replace("Z", "+00:00")).strftime("%B %d, %Y at %I:%M %p (UTC)") ) final_df = final_df.dropna() return final_df,order_df # Async display function async def display_results(): both_dfs = await main() final_df = both_dfs[0] order_df = both_dfs[1] dollar_delta = final_df['dollar delta'].sum() token_delta = final_df['token delta'].sum() avg_position = abs(dollar_delta / token_delta) metrics = ( f"**Dollar delta:** {dollar_delta:.2f} \n" f"**Token delta:** {token_delta:.2f} \n" f"**Avg position:** {avg_position:.6f}" ) return final_df,order_df,metrics # Gradio UI with async load with gr.Blocks() as demo: gr.Markdown("# Patchy Trades") df_output = gr.Dataframe(label="All Trades") order_output = gr.Dataframe(label="All Jupiter Orders") metrics_output = gr.Markdown(label="Metrics Summary") demo.load(fn=display_results, outputs=[df_output, order_output,metrics_output]) demo.launch()