""" Production trading node (v2) — integrates all Phase 1-4 modules. Replaces live/node.py with modular architecture: - Data: HyperliquidDataProvider + HyperliquidCollector - Analytics: AnalyticsPipeline (book → OBI, VPIN, microprice, signals) - Risk: Treasury (positions, PnL, circuit breakers) - Filter: ToxicityFilter (pre-trade VPIN gating) - Maker: HlMakerPool (A-S quoting per coin) - Monitors: CrossVenueMonitor, FundingBasisMonitor, LiquidationRiskOverlay - Dashboard: writes metrics to JSON for dashboard server Usage: python -m live.node_v2 --testnet --coins BTC,ETH --mode paper """ from __future__ import annotations import asyncio import json import logging import os import sys import time from pathlib import Path from typing import Optional sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) from sim.maker import AvellanedaStoikovMaker, MakerConfig, Quote from sim.queue import QueueModel from live.filters.toxicity import ToxicityFilter from live.treasury import Treasury from live.integrator import AnalyticsPipeline from live.monitors.cross_venue import CrossVenueMonitor from live.monitors.funding_basis import FundingBasisMonitor from live.monitors.liq_risk import LiquidationRiskOverlay from sim.fills import QueueAwareFillModel from live.makers.hl_btc_eth import HlMakerPool logger = logging.getLogger("ftdt-node-v2") TESTNET_API = "https://api.hyperliquid-testnet.xyz/info" MAINNET_API = "https://api.hyperliquid.xyz/info" DEFAULT_COINS = ["BTC", "ETH"] class ProductionNode: """Production trading node integrating analytics, risk, maker, and monitors. Lifecycle per tick: 1. Fetch order books and mark prices from HL REST 2. Feed book/trade data into AnalyticsPipeline 3. Check Treasury circuit breakers 4. Check ToxicityFilter 5. Generate quotes via HlMakerPool 6. Place orders (paper or live) 7. Process fills, update Treasury 8. Write dashboard metrics 9. Check monitors (liq risk, funding, cross-venue) """ def __init__( self, coins: list[str] | None = None, testnet: bool = True, mode: str = "paper", # "paper" or "live" api_url: str | None = None, private_key: str | None = None, max_position_per_coin: float = 0.003, base_quote_size: float = 0.0002, initial_equity: float = 10000.0, tick_interval_sec: float = 2.0, metrics_file: str = "/tmp/ftdt-metrics-v2.json", ): self._coins = coins or DEFAULT_COINS self._testnet = testnet self._mode = mode self._api_url = api_url or (TESTNET_API if testnet else MAINNET_API) self._pk = private_key self._tick_interval = tick_interval_sec self._metrics_file = metrics_file # Core modules self._treasury = Treasury( initial_equity=initial_equity, max_position_per_asset=max_position_per_coin, ) self._pipelines = { coin: AnalyticsPipeline() for coin in coins } # Maker pool self._maker_pool = HlMakerPool( treasury=self._treasury, maker_config={ "base_size": base_quote_size, "max_spread_bps": 15.0, "vpin_threshold": 0.3, "vpin_alarm": 0.5, }, ) for coin in coins: self._maker_pool.add_maker(coin.upper(), max_inventory=max_position_per_coin) # Queue-aware fill model for realistic paper trading self._fill_model = QueueAwareFillModel() # WQI predictor — directional strategy from queue imbalance from strategies.wqi_predictor import WQIPredictor self._wqi_predictors = { coin: WQIPredictor(z_entry=2.0, max_hold_seconds=30, stop_loss_bps=5.0, take_profit_bps=10.0, size=base_quote_size, fee_model="taker") for coin in coins } # Funding arb strategy from strategies.funding_arb_strategy import FundingArb self._funding_arb = FundingArb( apr_threshold=0.30, apr_exit=0.10, size=base_quote_size * 5, max_hold_hours=48.0, taker_fee_pct=0.00045, ) # Monitors self._cross_venue = CrossVenueMonitor() self._funding_monitor = FundingBasisMonitor() # State self._running = False self._tick = 0 self._equity_history: list[dict] = [] async def start(self): logger.info("Node v2 starting — %d coins, mode=%s, testnet=%s", len(self._coins), self._mode, self._testnet) self._running = True self._equity_history.append({"t": time.time(), "v": self._treasury.equity}) async def stop(self): self._running = False logger.info("Node v2 stopped — PnL: $%.2f (%.2f%%), %d trades", self._treasury.total_pnl(), self._treasury.pnl_pct(), self._treasury._daily_trades) async def run(self): """Main event loop.""" await self.start() try: while self._running: try: await self._tick_cycle() except Exception as e: logger.error("Tick error: %s", e, exc_info=True) self._treasury.record_api_error() await asyncio.sleep(self._tick_interval) finally: await self.stop() async def _tick_cycle(self): self._tick += 1 # 1. Fetch data prices = await self._fetch_mark_prices() books = {} for coin in self._coins: book = await self._fetch_orderbook(coin) if book: books[coin] = book prices[coin] = book.get("mid", prices.get(coin, 0)) # 2. Feed analytics pipeline for coin in self._coins: book = books.get(coin, {}) pipeline = self._pipelines[coin] if book.get("bids") and book.get("asks"): bids = {float(px): float(sz) for px, sz in book.get("bids", [])} asks = {float(px): float(sz) for px, sz in book.get("asks", [])} pipeline.update_book(bids, asks) # 3. Check circuit breakers if self._treasury.is_halted(): if self._tick % 30 == 0: logger.warning("Circuit breaker halted: %s", self._treasury.halt_reason) self._write_metrics() return # 4. Update makers with prices mid_prices = {coin: self._pipelines[coin].mid for coin in self._coins} self._maker_pool.observe_all(mid_prices) # 5. Update book info on makers for coin in self._coins: maker = self._maker_pool.get(coin) book = books.get(coin, {}) if maker and book: bids = dict(book.get("bids", []) or []) asks = dict(book.get("asks", []) or []) bb = max(bids) if bids else 0 ba = min(asks) if asks else 0 maker.update_book(bb, ba) # 6. Generate quotes quotes = self._maker_pool.quote_all() # 7. Simulate fills (paper mode — queue-aware) if self._mode == "paper": for coin in self._coins: q = quotes.get(coin) pipeline = self._pipelines[coin] if q: self._simulate_paper_fills(coin, q, pipeline) if self._tick % 3 == 0: self._simulate_wqi_trades(coin, pipeline) # 8. Update funding monitor and check funding arb for coin in self._coins: funding = await self._fetch_funding(coin) if funding is not None: self._funding_monitor.update_funding(coin, funding) # 9. Update cross-venue for coin in self._coins: self._cross_venue.update("hl", coin, mid_prices.get(coin, 0), time.time()) # 10. Check liquidation risk liq_overlay = LiquidationRiskOverlay(self._treasury) for coin, status in liq_overlay.check_all().items(): if status["level"] in ("danger", "critical"): logger.warning("Liquidation risk [%s]: %s — distance %.1f%%", coin, status["level"], status["distance_pct"]) # 11. Write metrics self._write_metrics() # 12. Log summary if self._tick % 30 == 0: self._log_status() # ── Data fetching ───────────────────────────────────────── async def _fetch_mark_prices(self) -> dict[str, float]: import requests try: resp = requests.post(self._api_url, json={"type": "metaAndAssetCtxs"}, timeout=10) data = resp.json() if isinstance(data, list) and len(data) >= 2: universe = data[0].get("universe", []) ctxs = data[1] prices = {} for i, u in enumerate(universe): name = u.get("name", "") if name in self._coins and i < len(ctxs): prices[name] = float(ctxs[i].get("markPx", 0)) return prices except Exception as e: logger.debug("Mark price fetch error: %s", e) return {} async def _fetch_orderbook(self, coin: str) -> dict | None: import requests try: resp = requests.post(self._api_url, json={"type": "l2Book", "coin": coin}, timeout=5) data = resp.json() levels = data.get("levels", []) if levels and len(levels) >= 2: bids = [(float(l["px"]), float(l["sz"])) for l in levels[0] if float(l["sz"]) > 0] asks = [(float(l["px"]), float(l["sz"])) for l in levels[1] if float(l["sz"]) > 0] bb = bids[0][0] if bids else 0 ba = asks[0][0] if asks else 0 return {"bids": bids, "asks": asks, "mid": (bb + ba) / 2 if bb and ba else 0} except Exception as e: logger.debug("Orderbook fetch error for %s: %s", coin, e) return None async def _fetch_funding(self, coin: str) -> float | None: import requests try: resp = requests.post(self._api_url, json={"type": "metaAndAssetCtxs"}, timeout=5) data = resp.json() if isinstance(data, list) and len(data) >= 2: universe = data[0].get("universe", []) ctxs = data[1] for i, u in enumerate(universe): if u.get("name", "") == coin and i < len(ctxs): return float(ctxs[i].get("funding", "0")) except Exception: pass return None # ── Paper trading ──────────────────────────────────────── def _simulate_paper_fills(self, coin: str, quote, pipeline: AnalyticsPipeline): mid = pipeline.mid if mid <= 0 or quote is None: return bid_fill = self._fill_model.check_fill( aggressor_side="sell", agg_size=pipeline._depth_ask or 0.1, agg_price=max(getattr(quote, "bid", mid) - 1, 1), our_price=getattr(quote, "bid", mid), our_size=getattr(quote, "bid_size", 0.0002), depth_ahead=self._fill_model.estimate_depth_ahead( our_price=getattr(quote, "bid", mid), our_side="bid", best_bid=pipeline._best_bid, best_ask=pipeline._best_ask, bid_depth=pipeline._depth_bid or 1.0, ask_depth=pipeline._depth_ask or 1.0, ), ) if bid_fill["filled"]: side = "buy" size = bid_fill["fill_size"] px = getattr(quote, "bid", mid) fee = size * px * 0.0002 can = self._treasury.can_open(coin, side, size, px) if can["allowed"]: self._treasury.record_fill(coin, side, size, px, fee, pnl=0) maker = self._maker_pool.get(coin) if maker: maker.record_fill(side, size, px, fee) ask_fill = self._fill_model.check_fill( aggressor_side="buy", agg_size=pipeline._depth_bid or 0.1, agg_price=min(getattr(quote, "ask", mid) + 1, mid * 2), our_price=getattr(quote, "ask", mid), our_size=getattr(quote, "ask_size", 0.0002), depth_ahead=self._fill_model.estimate_depth_ahead( our_price=getattr(quote, "ask", mid), our_side="ask", best_bid=pipeline._best_bid, best_ask=pipeline._best_ask, bid_depth=pipeline._depth_bid or 1.0, ask_depth=pipeline._depth_ask or 1.0, ), ) if ask_fill["filled"]: side = "sell" size = ask_fill["fill_size"] px = getattr(quote, "ask", mid) fee = size * px * 0.0002 can = self._treasury.can_open(coin, side, size, px) if can["allowed"]: self._treasury.record_fill(coin, side, size, px, fee, pnl=0) maker = self._maker_pool.get(coin) if maker: maker.record_fill(side, size, px, fee) def _simulate_wqi_trades(self, coin: str, pipeline: AnalyticsPipeline): mid = pipeline.mid if mid <= 0: return predictor = self._wqi_predictors.get(coin) if predictor is None: return bids_list = [(pipeline._best_bid, pipeline._depth_bid or 1.0)] asks_list = [(pipeline._best_ask, pipeline._depth_ask or 1.0)] signal = predictor.feed_signal(bids_list, asks_list, mid) if signal["action"] in ("BUY", "SELL"): side = signal["action"].lower() size = 0.0002 px = mid fee = size * px * 0.0005 can = self._treasury.can_open(coin, side, size, px) if can["allowed"]: self._treasury.record_fill(coin, side, size, px, fee, pnl=0) fee_paid = size * px * 0.0005 self._treasury._fees_paid += fee_paid logger.info(f"[WQI-{coin}] {signal['action']} signal: " f"z={signal['z_score']:.2f} wqi={signal['wqi']:.3f} " f"reason={signal['reason']}") elif signal["action"] == "EXIT": pos = self._treasury.position(coin) if abs(pos) > 0: side = "sell" if pos > 0 else "buy" fee = abs(pos) * mid * 0.0005 pnl = pos * (mid - predictor._entry_price) if predictor._entry_price > 0 else 0 self._treasury.record_fill(coin, side, abs(pos), mid, fee, pnl) logger.info(f"[WQI-{coin}] EXIT: z={signal['z_score']:.2f} " f"pnl=${pnl:.4f} reason={signal['reason']}") # ── Dashboard ──────────────────────────────────────────── def _write_metrics(self): equity = self._treasury.equity t = time.time() self._equity_history.append({"t": t, "v": equity}) if len(self._equity_history) > 600: self._equity_history = self._equity_history[-600:] try: data = { "timestamp": t, "treasury": self._treasury.summary(), "analytics": {c: p.emit() for c, p in self._pipelines.items()}, "maker": self._maker_pool.summary(), "wqi": {c: p.summary() for c, p in self._wqi_predictors.items()}, "funding_arb": self._funding_arb.summary(), "fill_model": { "fill_rate": round(self._fill_model.fill_rate(), 4), "fills": self._fill_model.fill_count, "skips": self._fill_model.skip_count, }, "funding": self._funding_monitor.summary(), "cross_venue": self._cross_venue.summary("BTC"), "equity_history": self._equity_history, } with open(self._metrics_file, "w") as f: json.dump(data, f, default=str) except IOError: pass def _log_status(self): treasury = self._treasury.summary() logger.info( "Tick %d | Equity: $%.0f | PnL: %.2f%% | Trades: %d | Positions: %s", self._tick, treasury["equity"], treasury["pnl_pct"], treasury["daily_trades"], treasury["positions"], ) # ── CLI ───────────────────────────────────────────────────── async def _main(): import argparse p = argparse.ArgumentParser(description="FTDT Quant Lab — Production Node v2") p.add_argument("--coins", default="BTC,ETH", help="Comma-separated coin list") p.add_argument("--testnet", action="store_true", default=True) p.add_argument("--mainnet", dest="testnet", action="store_false") p.add_argument("--mode", default="paper", choices=["paper", "live"]) p.add_argument("--max-position", type=float, default=0.003) p.add_argument("--base-size", type=float, default=0.0002) p.add_argument("--equity", type=float, default=10000.0) p.add_argument("--tick-interval", type=float, default=2.0) p.add_argument("--metrics-file", default="/tmp/ftdt-metrics-v2.json") p.add_argument("--private-key", default=None) args = p.parse_args() logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(name)s] %(message)s", datefmt="%H:%M:%S", ) coins = [c.strip().upper() for c in args.coins.split(",")] node = ProductionNode( coins=coins, testnet=args.testnet, mode=args.mode, private_key=args.private_key, max_position_per_coin=args.max_position, base_quote_size=args.base_size, initial_equity=args.equity, tick_interval_sec=args.tick_interval, metrics_file=args.metrics_file, ) try: await node.run() except KeyboardInterrupt: logger.info("Shutting down...") if __name__ == "__main__": asyncio.run(_main())