00bb06434e
data/collectors/hyperliquid.py: - Fixed _handle_trades to handle HL WebSocket trade data as list[dict] (each trade dict has its own 'coin' field) instead of assuming single dict - Fixed _poll_funding and _poll_open_interest to guard against data[0] being list vs dict (mainnet vs testnet response format difference) data/duckdb_load.py: - Replaced deprecated duckdb.from_sequence() with executemany() for l2_snapshots, trades, and funding table inserts (DuckDB 1.5.x API) Tick runner verified on real mainnet BTC data: - 324 events (45 L2 + 279 trades), VPIN-gated A-S MM, 1 trade, 98 cancels - HFT dashboard generated with 8 panels from real tick data
467 lines
19 KiB
Python
467 lines
19 KiB
Python
"""
|
|
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 | list, local_ts: float):
|
|
trade_list = data if isinstance(data, list) else [data]
|
|
for trade in trade_list:
|
|
coin = trade.get("coin", "?") if isinstance(trade, dict) else "?"
|
|
if coin not in self._coins:
|
|
continue
|
|
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:
|
|
raw_universe = data[0]
|
|
if isinstance(raw_universe, dict):
|
|
universe = raw_universe.get("universe", [])
|
|
else:
|
|
universe = raw_universe if isinstance(raw_universe, list) else []
|
|
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:
|
|
raw_universe = data[0]
|
|
if isinstance(raw_universe, dict):
|
|
universe = raw_universe.get("universe", [])
|
|
else:
|
|
universe = raw_universe if isinstance(raw_universe, list) else []
|
|
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())
|