Files
ftdt-quant-lab/live/node_v2.py
T
ramseshk 3073415d33 feat: funding arb strategy, queue-aware paper fills, WQI live integration
- strategies/funding_arb_strategy.py: full backtestable funding rate carry module
  with entry/exit thresholds, position tracking, funding payment accounting,
  basis stop-loss, max-hold timeout. Includes backtest_funding_arb() and
  run_funding_discovery() for threshold optimization
- live/node_v2.py: replaced naive random fills with QueueAwareFillModel (sim/fills.py)
  with queue-priority simulation; integrated WQI predictor and funding arb strategies;
  per-coin WQI signal generation every 3 ticks; funding arb metrics in dashboard
- cli.py: added 'funding' command for funding rate distribution analysis and
  threshold backtesting
- tests/test_funding_arb.py: 20 tests covering entry/exit logic, fee accounting,
  signal generation, backtesting, and node integration

321 tests passing (20 new).
2026-08-11 11:15:25 +08:00

478 lines
18 KiB
Python

"""
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())