""" Hyperliquid real-time data collector. Streams: L2 order book (full book maintained per coin), trades, mark prices (allMids), and user notifications (fills, liquidations). Polls: funding rates, open interest, predicted funding (REST). All raw messages are stored via RawMessageStore. Order books are maintained with full incremental reconstruction + gap detection. Usage: collector = HyperliquidCollector( store=store, coins=["BTC", "ETH"], testnet=True, ) await collector.start() await collector.run() # blocks until Ctrl+C """ from __future__ import annotations import asyncio import json import logging import time from datetime import datetime, timezone from typing import Optional from data.store import RawMessageStore from data.normalizer import normalize_timestamp, SequenceTracker from data.latency import LatencyTracker logger = logging.getLogger(__name__) TESTNET_WS = "wss://api.hyperliquid-testnet.xyz/ws" MAINNET_WS = "wss://api.hyperliquid.xyz/ws" TESTNET_API = "https://api.hyperliquid-testnet.xyz/info" MAINNET_API = "https://api.hyperliquid.xyz/info" LEDGER_DECIMALS = {"BTC": 5, "ETH": 6, "SOL": 7, "HYPE": 6, "VVV": 6} class OrderBook: """Reconstructed limit order book for one coin.""" def __init__(self, coin: str): self.coin = coin self.bids: dict[float, float] = {} # price → size self.asks: dict[float, float] = {} self._seq: int = 0 self._update_count: int = 0 self._snapshot_count: int = 0 def apply_snapshot(self, levels: list, side: str): """Full replace of one side.""" target = self.bids if side == "bids" else self.asks target.clear() for level in levels: px = float(level["px"]) sz = float(level["sz"]) if sz > 0: target[px] = sz if side == "bids": self._snapshot_count += 1 def apply_update(self, delta: dict): """Apply incremental update to one side.""" side = "bids" if delta.get("side") == "B" else "asks" target = self.bids if side == "bids" else self.asks px = float(delta["px"]) sz = float(delta["sz"]) if sz == 0: target.pop(px, None) else: target[px] = sz self._update_count += 1 def best_bid(self) -> float: return max(self.bids) if self.bids else 0.0 def best_ask(self) -> float: return min(self.asks) if self.asks else 0.0 def mid(self) -> float: bb = self.best_bid() ba = self.best_ask() return (bb + ba) / 2.0 if bb and ba else 0.0 def total_depth(self, side: str, levels: int = 10) -> float: target = self.bids if side == "bids" else self.asks return sum(sorted(target.values(), reverse=(side == "bids"))[:levels]) def stats(self) -> dict: bb = self.best_bid() ba = self.best_ask() spread = ba - bb if bb and ba else 0 return { "coin": self.coin, "best_bid": bb, "best_ask": ba, "mid": (bb + ba) / 2.0 if bb and ba else 0.0, "spread": spread, "spread_bps": round((spread / bb * 10000), 1) if bb else 0, "bid_levels": len(self.bids), "ask_levels": len(self.asks), "bid_depth_10": self.total_depth("bids", 10), "ask_depth_10": self.total_depth("asks", 10), "snapshots": self._snapshot_count, "updates": self._update_count, "seq": self._seq, } class HyperliquidCollector: """Streams and stores Hyperliquid market data.""" def __init__( self, store: RawMessageStore, coins: list[str] | None = None, testnet: bool = True, poll_interval_sec: float = 60.0, reconnect_delay: float = 2.0, ): self._store = store self._coins = coins or ["BTC", "ETH"] self._testnet = testnet self._poll_interval = poll_interval_sec self._reconnect_delay = reconnect_delay self._ws_url = TESTNET_WS if testnet else MAINNET_WS self._api_url = TESTNET_API if testnet else MAINNET_API self._books: dict[str, OrderBook] = {c: OrderBook(c) for c in self._coins} self._seq_tracker = SequenceTracker() self._latency = LatencyTracker() self._running = False @property def books(self) -> dict[str, OrderBook]: return self._books @property def latency(self) -> LatencyTracker: return self._latency async def start(self): """Start the store and prepare.""" self._store.start() self._running = True logger.info("HyperliquidCollector started (%s, %d coins, %s)", "testnet" if self._testnet else "mainnet", len(self._coins), self._coins) async def stop(self): """Graceful shutdown.""" self._running = False self._store.stop() logger.info("HyperliquidCollector stopped") async def run(self): """Main entrypoint — blocks with WebSocket + REST pollers.""" await self.start() try: async with asyncio.TaskGroup() as tg: tg.create_task(self._ws_loop()) for task in [self._poll_funding, self._poll_open_interest, self._poll_liquidations, self._stats_reporter]: tg.create_task(task()) except ExceptionGroup as eg: for exc in eg.exceptions: logger.error("Collector error: %s", exc) finally: await self.stop() # ── WebSocket stream ───────────────────────────────────── async def _ws_loop(self): while self._running: try: await self._connect_and_stream() except Exception as e: logger.warning("WebSocket error: %s — reconnecting in %.1fs", e, self._reconnect_delay) await asyncio.sleep(self._reconnect_delay) async def _connect_and_stream(self): try: import websockets except ImportError: logger.error("websockets not installed; pip install websockets") return async with websockets.connect(self._ws_url, ping_interval=30, ping_timeout=10) as ws: for coin in self._coins: await ws.send(json.dumps({"method": "subscribe", "subscription": {"type": "l2Book", "coin": coin}})) await ws.send(json.dumps({"method": "subscribe", "subscription": {"type": "trades", "coin": coin}})) await ws.send(json.dumps({"method": "subscribe", "subscription": {"type": "allMids"}})) logger.info("Subscribed to %d coins (l2Book, trades, allMids)", len(self._coins)) while self._running: try: raw = await asyncio.wait_for(ws.recv(), timeout=30) except asyncio.TimeoutError: continue local_ts = time.time() try: msg = json.loads(raw) except json.JSONDecodeError: continue channel = msg.get("channel", "") data = msg.get("data", {}) if channel == "l2Book": await self._handle_l2book(data, local_ts) elif channel == "trades": await self._handle_trades(data, local_ts) elif channel == "allMids": await self._handle_all_mids(data, local_ts) elif channel == "subscriptionResponse": logger.debug("Subscription confirmed: %s", data) async def _handle_l2book(self, data: dict, local_ts: float): coin = data.get("coin", "?") if coin not in self._coins: return levels = data.get("levels", []) time_ms = normalize_timestamp(data.get("time"), source="hl") if levels and isinstance(levels[0], list): # Snapshot: levels = [[bids...], [asks...]] book = self._books[coin] bid_levels = [] ask_levels = [] for bid in levels[0]: if float(bid.get("sz", 0)) > 0: bid_levels.append({"px": bid["px"], "sz": bid["sz"]}) for ask in levels[1]: if float(ask.get("sz", 0)) > 0: ask_levels.append({"px": ask["px"], "sz": ask["sz"]}) book.apply_snapshot(bid_levels, "bids") book.apply_snapshot(ask_levels, "asks") self._seq_tracker.reset("l2book", coin) else: # Incremental update side = "bids" if data.get("side") == "B" else "asks" delta = {"side": data.get("side", "B"), "px": data["px"], "sz": data["sz"]} self._books[coin].apply_update(delta) seq = time_ms gap = self._seq_tracker.check("l2book", coin, seq) if gap: logger.warning("L2 gap %s: expected=%d got=%d gap=%d", coin, gap["expected"], gap["got"], gap["gap_size"]) self._store.push( channel="l2book", coin=coin, exchange_ts=time_ms, payload={"type": "snapshot" if levels and isinstance(levels[0], list) else "delta", "levels": levels if isinstance(levels, list) else {}, "delta": {} if levels and isinstance(levels[0], list) else { "side": data.get("side", ""), "px": data.get("px", ""), "sz": data.get("sz", ""), }}, ) self._latency.record_transport(time_ms, local_ts) async def _handle_trades(self, data: dict, local_ts: float): coin = data.get("coin", "?") if coin not in self._coins: return trade_list = data if isinstance(data, list) else [data] for trade in trade_list: time_ms = normalize_timestamp(trade.get("time"), source="hl") self._store.push( channel="trades", coin=coin, exchange_ts=time_ms, payload={ "side": trade.get("side", ""), "px": trade.get("px", ""), "sz": trade.get("sz", ""), "hash": trade.get("hash", ""), }, ) self._latency.record_transport(time_ms, local_ts) async def _handle_all_mids(self, data: dict, local_ts: float): mids = data.get("mids", {}) time_ms = int(data.get("time", time.time() * 1000)) for asset, mid_px in mids.items(): if asset in self._coins: self._store.push( channel="mark", coin=asset, exchange_ts=time_ms, payload={"mark_px": mid_px}, ) self._latency.record_transport(time_ms, local_ts) # ── REST pollers ───────────────────────────────────────── async def _poll_funding(self): import aiohttp while self._running: try: async with aiohttp.ClientSession() as session: async with session.post(self._api_url, json={"type": "metaAndAssetCtxs"}, timeout=aiohttp.ClientTimeout(total=10)) as resp: data = await resp.json() if isinstance(data, list) and len(data) >= 2: universe = data[0].get("universe", []) ctxs = data[1] now_ms = int(time.time() * 1000) for i, asset_info in enumerate(universe): name = asset_info.get("name", "") if name not in self._coins or i >= len(ctxs): continue ctx = ctxs[i] self._store.push( channel="funding", coin=name, exchange_ts=now_ms, payload={ "funding": ctx.get("funding", "0"), "mark_px": ctx.get("markPx", "0"), "index_px": ctx.get("oraclePx", ctx.get("indexPx", "0")), "open_interest": ctx.get("openInterest", "0"), "day_ntl_volume": ctx.get("dayNtlVlm", "0"), }, ) # Predicted funding async with session.post(self._api_url, json={"type": "predictedFundings"}, timeout=aiohttp.ClientTimeout(total=10)) as resp: pred = await resp.json() if isinstance(pred, list): now_ms = int(time.time() * 1000) for item in pred: name = item.get("name", "") if name in self._coins: self._store.push( channel="predicted_funding", coin=name, exchange_ts=now_ms, payload={"funding": item.get("funding", "0"), "premium": item.get("premium", "0")}, ) except Exception as e: logger.warning("Funding poll error: %s", e) await asyncio.sleep(self._poll_interval) async def _poll_open_interest(self): import aiohttp while self._running: try: async with aiohttp.ClientSession() as session: async with session.post(self._api_url, json={"type": "metaAndAssetCtxs"}, timeout=aiohttp.ClientTimeout(total=10)) as resp: data = await resp.json() if isinstance(data, list) and len(data) >= 2: universe = data[0].get("universe", []) ctxs = data[1] now_ms = int(time.time() * 1000) for i, asset_info in enumerate(universe): name = asset_info.get("name", "") if name not in self._coins or i >= len(ctxs): continue oi = ctxs[i].get("openInterest", "0") self._store.push( channel="open_interest", coin=name, exchange_ts=now_ms, payload={"open_interest": oi}, ) except Exception as e: logger.warning("OI poll error: %s", e) await asyncio.sleep(self._poll_interval) async def _poll_liquidations(self): """Poll for recent liquidation events (public feed approximation). Hyperliquid doesn't have a public liquidation-only endpoint, so we poll trade history and filter for liquidations. This is a best-effort approximation — full liquidation data requires processing the trades feed in real-time and checking the 'liquidation' field.""" import aiohttp while self._running: try: for coin in self._coins: now = int(time.time() * 1000) async with aiohttp.ClientSession() as session: async with session.post( self._api_url, json={ "type": "userFillsByTime", "user": "0x0000000000000000000000000000000000000000", "startTime": now - 3600_000, "limit": 200, }, timeout=aiohttp.ClientTimeout(total=10), ) as resp: fills = await resp.json() if isinstance(fills, list): for fill in fills: if fill.get("liquidation") and fill.get("coin", "").upper() in self._coins: time_ms = normalize_timestamp(fill.get("time"), source="hl") self._store.push( channel="liquidation", coin=fill["coin"].upper(), exchange_ts=time_ms, payload={ "side": fill.get("side", ""), "sz": fill.get("sz", ""), "px": fill.get("px", ""), }, ) except Exception: pass await asyncio.sleep(self._poll_interval * 5) async def _stats_reporter(self): while self._running: await asyncio.sleep(60) for book in self._books.values(): logger.info("Book %s: %s", book.coin, book.stats()) logger.info("Latency: %s", self._latency.summary()) logger.info("Store: %d total written", self._store.total_written) # ── CLI entrypoint ─────────────────────────────────────────── async def _main(): import argparse p = argparse.ArgumentParser(description="Hyperliquid data collector") p.add_argument("--coins", nargs="+", default=["BTC", "ETH"]) p.add_argument("--testnet", action="store_true", default=True) p.add_argument("--mainnet", dest="mainnet", action="store_true") p.add_argument("--data-dir", default="data/raw") p.add_argument("--poll-interval", type=float, default=60.0) p.add_argument("--flush-interval", type=float, default=5.0) args = p.parse_args() logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(name)s] %(message)s", datefmt="%H:%M:%S") store = RawMessageStore(data_dir=args.data_dir, flush_interval_sec=args.flush_interval) collector = HyperliquidCollector( store=store, coins=args.coins, testnet=not args.mainnet, poll_interval_sec=args.poll_interval, ) try: await collector.run() except KeyboardInterrupt: logger.info("Shutting down...") if __name__ == "__main__": asyncio.run(_main())