feat: Phase 5 — integration layer (analytics pipeline, production node v2, CLI) + 12 tests

live/integrator.py (AnalyticsPipeline):
  Real-time pipeline: data → microstructure → signals.
  Accumulates book snapshots + trades, computes OBI, VPIN, microprice,
  spread, depth, trade imbalance, HFT regime, and emits composite
  signal with confidence and breakdown. Per-coin isolation.

live/node_v2.py (ProductionNode):
  Rebuilt production node integrating ALL Phase 1-4 modules:
  - REST data fetching (order book, mark prices, funding rates)
  - AnalyticsPipeline per coin for real-time microstructure signals
  - Treasury for position/capital/PnL/breaker management
  - ToxicityFilter integration via HlMakerPool makers
  - HlMakerPool for per-coin A-S quoting
  - CrossVenueMonitor, FundingBasisMonitor, LiquidationRiskOverlay
  - Paper trading with probabilistic fill simulation
  - Dashboard metrics JSON output (equity, treasury, analytics, maker)
  - Periodic status logging

cli.py (unified CLI):
  Subcommands integrating all modules:
    collect   — Run Hyperliquid data collector to Parquet
    analyze   — Run microstructure analytics on stored data
    simulate  — Run market-making simulator on stored data
    run       — Start production trading node (paper or live)
    backtest  — Run VectorBT backtest

12 integration tests (all pass):
  - AnalyticsPipeline: empty, book, trade, VPIN, emit, regime, isolation
  - ProductionNode: creation, tick cycle (3 ticks), metrics JSON output
  - CLI: import verification

Total test suite: 184 tests, all passing.
This commit is contained in:
ramseshk
2026-08-07 14:54:47 +08:00
parent 4f66ef36a9
commit e58c5951b7
4 changed files with 1042 additions and 0 deletions
+314
View File
@@ -0,0 +1,314 @@
"""
FTDT Quant Lab — unified CLI.
Subcommands:
collect — Run data collector (streams to Parquet)
analyze — Run analytics on stored data (Phase 2)
simulate — Run market-making simulator on stored data (Phase 3)
run — Start production trading node (Phase 4)
backtest — Run VectorBT backtest (existing)
Usage:
python -m cli collect --coins BTC,ETH --data-dir data/raw
python -m cli analyze --data-dir data/raw --start 2026-08-01 --end 2026-08-07
python -m cli simulate --data-dir data/raw --coin BTC --hours 24
python -m cli run --coins BTC,ETH --mode paper
python -m cli backtest --strategy pairs --interval 1h
"""
from __future__ import annotations
import asyncio
import logging
import os
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent))
def cmd_collect(args):
"""Run the Hyperliquid data collector."""
from data.collectors.hyperliquid import HyperliquidCollector
from data.store import RawMessageStore
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,
)
asyncio.run(collector.run())
def cmd_analyze(args):
"""Run microstructure analytics on stored data."""
from data.store import read_range
print(f"Reading {args.channel}/{args.coin} from {args.start_date} to {args.end_date}...")
messages = read_range(
args.data_dir,
channel=args.channel,
coin=args.coin.upper(),
start_date=args.start_date,
end_date=args.end_date,
)
print(f"Loaded {len(messages)} messages")
if args.channel == "l2book":
from microstructure.book import batch_book_stats
snapshots = []
for msg in messages:
payload = msg["payload"]
levels = payload.get("levels", [])
if levels and isinstance(levels, list) and len(levels) >= 2:
bids = {}
asks = {}
for bid in levels[0]:
if float(bid.get("sz", 0)) > 0:
bids[float(bid["px"])] = float(bid["sz"])
for ask in levels[1]:
if float(ask.get("sz", 0)) > 0:
asks[float(ask["px"])] = float(ask["sz"])
snapshots.append({"bids": bids, "asks": asks})
stats = batch_book_stats(snapshots)
print(json.dumps(stats, indent=2, default=str))
elif args.channel == "trades":
from microstructure.trades import classify_bulk_lee_ready, trade_arrival_rate, trade_volume_profile
trades = [msg["payload"] for msg in messages]
mids = [float(msg["payload"].get("px", 0)) for msg in messages]
times = [msg["exchange_ts"] for msg in messages]
sides = classify_bulk_lee_ready(trades, mids)
buys = sum(1 for s in sides if s == "buy")
sells = sum(1 for s in sides if s == "sell")
arrival = trade_arrival_rate(times)
vol = trade_volume_profile(trades)
print(f"Trades: {len(trades)} total ({buys} buy, {sells} sell)")
print(f"Arrival rate: {json.dumps(arrival, indent=2, default=str)}")
print(f"Volume profile: {json.dumps(vol, indent=2, default=str)}")
elif args.channel == "funding":
from microstructure.funding import funding_regime, basis_spread
rates = [float(msg["payload"].get("funding", 0)) for msg in messages]
marks = [float(msg["payload"].get("mark_px", 0)) for msg in messages]
regime = funding_regime(rates, window_hours=24, n_samples_per_hour=1)
print(f"Funding regime: {json.dumps(regime, indent=2, default=str)}")
else:
print(f"Channel '{args.channel}' — raw dump:")
for msg in messages[:5]:
print(json.dumps(msg, indent=2, default=str))
if len(messages) > 5:
print(f"... and {len(messages) - 5} more")
def cmd_simulate(args):
"""Run market-making simulator on stored data with L2 events."""
from data.store import read_range
from sim.engine import SimulationEngine, SimConfig
from sim.maker import MakerConfig
print(f"Loading L2 book data for {args.coin} from {args.start_date} to {args.end_date}...")
l2_messages = read_range(
args.data_dir,
channel="l2book",
coin=args.coin.upper(),
start_date=args.start_date,
end_date=args.end_date,
)
print(f"Loaded {len(l2_messages)} L2 updates")
trade_messages = read_range(
args.data_dir,
channel="trades",
coin=args.coin.upper(),
start_date=args.start_date,
end_date=args.end_date,
)
print(f"Loaded {len(trade_messages)} trades")
events = []
for msg in l2_messages:
payload = msg["payload"]
levels = payload.get("levels", [])
bids = {}
asks = {}
if levels and isinstance(levels, list) and len(levels) >= 2:
for bid in levels[0]:
if float(bid.get("sz", 0)) > 0:
bids[float(bid["px"])] = float(bid["sz"])
for ask in levels[1]:
if float(ask.get("sz", 0)) > 0:
asks[float(ask["px"])] = float(ask["sz"])
events.append({
"type": "l2",
"data": {"bids": bids, "asks": asks},
"time": msg["local_ts"],
"coin": args.coin.upper(),
})
for msg in trade_messages:
payload = msg["payload"]
events.append({
"type": "trade",
"data": payload,
"time": msg["local_ts"],
"coin": args.coin.upper(),
})
events.sort(key=lambda e: e["time"])
print(f"Total events: {len(events)}")
config = SimConfig(
maker=MakerConfig(
base_size=args.base_size,
max_inventory=args.max_inventory,
gamma=args.gamma,
),
max_inventory=args.max_inventory,
cancel_after_ms=args.cancel_after_ms,
quote_refresh_ms=args.quote_refresh_ms,
seed=args.seed,
)
engine = SimulationEngine(config=config, seed=args.seed)
engine.run(events)
stats = engine.stats()
breakdown = engine.breakdown()
print("\n=== Simulation Results ===")
print(f"Duration: {events[-1]['time'] - events[0]['time']:.0f}s" if events else "0s")
print(f"Trades: {stats.total_trades} ({stats.bid_fills} bid, {stats.ask_fills} ask)")
print(f"Toxic fills: {stats.toxic_fills} ({stats.adverse_rate:.1%})")
print(f"Cancels: {stats.cancels}")
print(f"Avg spread: {stats.avg_spread_bps} bps")
print(f"Max inventory: {stats.max_inventory}")
print(f"Max drawdown: {stats.max_drawdown}%")
print(f"Sharpe: {stats.sharpe} Sortino: {stats.sortino}")
print(f"Uptime: {stats.uptime_pct}%")
print(f"\nPnL Breakdown:")
print(f" Spread capture: ${breakdown.spread_capture:.4f}")
print(f" Inventory PnL: ${breakdown.inventory_pnl:.4f}")
print(f" Maker fees: ${breakdown.maker_fees:.4f}")
print(f" Taker fees: ${breakdown.taker_fees:.4f}")
print(f" Funding PnL: ${breakdown.funding_pnl:.4f}")
print(f" Adverse selection: ${breakdown.adverse_selection_cost:.4f}")
print(f" ─────────────────────────────")
print(f" Gross PnL: ${breakdown.gross_pnl:.4f}")
print(f" Net PnL: ${breakdown.net_pnl:.4f}")
def cmd_run(args):
"""Start the production trading node."""
import asyncio
from live.node_v2 import ProductionNode
node = ProductionNode(
coins=args.coins,
testnet=not args.mainnet,
mode=args.mode,
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,
)
asyncio.run(node.run())
def cmd_backtest(args):
"""Run a VBT backtest (existing functionality)."""
from backtests.vbt_runner import VBTBacktestRunner
runner = VBTBacktestRunner()
result = runner.run_strategy(strategy=args.strategy, interval=args.interval, limit=args.limit)
import json as _json
print(_json.dumps({k: v for k, v in (result or {}).items()
if k not in ("trades", "equity_curve")}, indent=2, default=str))
if result and result.get("trades"):
print(f"\n{len(result['trades'])} trades")
def main():
import argparse
p = argparse.ArgumentParser(description="FTDT Quant Lab CLI")
sp = p.add_subparsers(dest="command", required=True)
# collect
pc = sp.add_parser("collect", help="Run data collector")
pc.add_argument("--coins", nargs="+", default=["BTC", "ETH"])
pc.add_argument("--mainnet", action="store_true")
pc.add_argument("--data-dir", default="data/raw")
pc.add_argument("--poll-interval", type=float, default=60.0)
pc.add_argument("--flush-interval", type=float, default=5.0)
# analyze
pa = sp.add_parser("analyze", help="Run microstructure analytics")
pa.add_argument("--data-dir", default="data/raw")
pa.add_argument("--channel", default="l2book", choices=["l2book", "trades", "funding", "mark", "open_interest", "liquidation"])
pa.add_argument("--coin", default="BTC")
pa.add_argument("--start-date", default="2026-08-01")
pa.add_argument("--end-date", default="2026-08-07")
# simulate
ps = sp.add_parser("simulate", help="Run market-making simulator")
ps.add_argument("--data-dir", default="data/raw")
ps.add_argument("--coin", default="BTC")
ps.add_argument("--start-date", default="2026-08-01")
ps.add_argument("--end-date", default="2026-08-07")
ps.add_argument("--gamma", type=float, default=0.1)
ps.add_argument("--base-size", type=float, default=0.001)
ps.add_argument("--max-inventory", type=float, default=0.005)
ps.add_argument("--cancel-after-ms", type=float, default=5000.0)
ps.add_argument("--quote-refresh-ms", type=float, default=2000.0)
ps.add_argument("--seed", type=int, default=42)
# run
pr = sp.add_parser("run", help="Start production node")
pr.add_argument("--coins", nargs="+", default=["BTC", "ETH"])
pr.add_argument("--mainnet", action="store_true")
pr.add_argument("--mode", default="paper", choices=["paper", "live"])
pr.add_argument("--max-position", type=float, default=0.003)
pr.add_argument("--base-size", type=float, default=0.0002)
pr.add_argument("--equity", type=float, default=10000.0)
pr.add_argument("--tick-interval", type=float, default=2.0)
pr.add_argument("--metrics-file", default="/tmp/ftdt-metrics-v2.json")
# backtest
pb = sp.add_parser("backtest", help="Run VBT backtest")
pb.add_argument("--strategy", default="pairs")
pb.add_argument("--interval", default="1h")
pb.add_argument("--limit", type=int, default=500)
args = p.parse_args()
import json as _json
import json
if args.command == "collect":
cmd_collect(args)
elif args.command == "analyze":
cmd_analyze(args)
elif args.command == "simulate":
cmd_simulate(args)
elif args.command == "run":
cmd_run(args)
elif args.command == "backtest":
cmd_backtest(args)
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s", datefmt="%H:%M:%S")
main()
+214
View File
@@ -0,0 +1,214 @@
"""
Real-time analytics pipeline: data → microstructure → signals.
Connects the data collector's output (order books, trades) to
microstructure analytics and produces actionable signals for
the maker pool and strategies.
Usage:
pipeline = AnalyticsPipeline()
pipeline.update_book(bids, asks)
pipeline.update_trade(px, sz, mid)
signals = pipeline.emit() # {obi, vpin, microprice, regime, composite, ...}
"""
from __future__ import annotations
from collections import deque
from typing import Optional
from microstructure.book import (
mid_price,
microprice,
order_book_imbalance,
spread_stats,
depth_resiliency,
)
from microstructure.trades import classify_lee_ready
from microstructure.toxicity import compute_vpin
from microstructure.signals import composite_signal, detect_hft_regime
class AnalyticsPipeline:
"""Real-time pipeline producing microstructure signals from book/trade data.
Maintains rolling windows of:
- Order book snapshots (for OBI, spread, depth)
- Trade volumes by side (for VPIN, trade imbalance)
- Mid prices (for volatility, markouts)
"""
def __init__(
self,
obi_window: int = 100,
vpin_window: int = 50,
vpin_bucket_size: float = 5.0,
trade_window: int = 500,
price_window: int = 300,
):
self._obi_window = obi_window
self._vpin_window = vpin_window
self._vpin_bucket_size = vpin_bucket_size
self._trade_window = trade_window
self._price_window = price_window
self._mid: float = 0.0
self._best_bid: float = 0.0
self._best_ask: float = 0.0
self._spread_bps: float = 0.0
self._microprice: float = 0.0
self._obi: float = 0.0
self._depth_bid: float = 0.0
self._depth_ask: float = 0.0
self._buy_vol: deque[float] = deque(maxlen=self._trade_window)
self._sell_vol: deque[float] = deque(maxlen=self._trade_window)
self._prices: deque[float] = deque(maxlen=self._price_window)
self._obis: deque[float] = deque(maxlen=self._obi_window)
self._trade_count: int = 0
self._current_vpin: float = 0.0
# ── Data ingestion ────────────────────────────────────────
def update_book(self, bids: dict[float, float], asks: dict[float, float]):
"""Feed an order book snapshot."""
if not bids or not asks:
return
bid_prices = sorted(bids.keys(), reverse=True)
ask_prices = sorted(asks.keys())
self._best_bid = bid_prices[0]
self._best_ask = ask_prices[0]
self._mid = (self._best_bid + self._best_ask) / 2.0
ss = spread_stats(bids, asks)
self._spread_bps = ss["spread_bps"]
self._microprice = microprice(bids, asks)
self._obi = order_book_imbalance(bids, asks)
self._obis.append(self._obi)
dr = depth_resiliency(bids, asks)
self._depth_bid = dr["bid_vol"]
self._depth_ask = dr["ask_vol"]
self._prices.append(self._mid)
def update_trade(self, px: float, sz: float, mid: float | None = None):
"""Feed a trade event."""
self._trade_count += 1
ref = mid if mid is not None else self._mid
side = classify_lee_ready(px, ref)
if side == "buy":
self._buy_vol.append(sz)
elif side == "sell":
self._sell_vol.append(sz)
self._recompute_vpin()
# ── Analytics computation ─────────────────────────────────
def _recompute_vpin(self):
result = compute_vpin(
list(self._buy_vol),
list(self._sell_vol),
volume_bucket_size=self._vpin_bucket_size,
n_buckets=self._vpin_window,
)
self._current_vpin = result.get("vpin_value", 0.0)
def ema_obi(self, alpha: float = 0.1) -> float:
"""Exponential moving average of OBI."""
vals = list(self._obis)
if not vals:
return 0.0
ema = vals[0]
for v in vals[1:]:
ema = alpha * v + (1 - alpha) * ema
return round(ema, 4)
def trade_imbalance(self, window: int | None = None) -> float:
"""Recent trade volume skew [-1, 1]."""
w = window or self._trade_window
bv = list(self._buy_vol)[-w:]
sv = list(self._sell_vol)[-w:]
total = sum(bv) + sum(sv)
return (sum(bv) - sum(sv)) / total if total > 0 else 0.0
def obi_volatility(self) -> float:
import math
vals = list(self._obis)
if len(vals) < 2:
return 0.0
mean = sum(vals) / len(vals)
return (sum((v - mean) ** 2 for v in vals) / len(vals)) ** 0.5
def spread_mean(self) -> float:
return self._spread_bps
def trade_rate(self, window_seconds: float = 60.0) -> float:
if self._trade_count == 0:
return 0.0
return self._trade_count / max(window_seconds, 1)
def hft_regime(self) -> str:
return detect_hft_regime(
obi_std=self.obi_volatility(),
spread_mean_bps=self.spread_mean(),
trade_rate_per_sec=self.trade_rate(),
vpin=self._current_vpin,
)
# ── Emit ──────────────────────────────────────────────────
def emit(self, funding_regime: str = "neutral") -> dict:
"""Produce a full signal report from accumulated data."""
self._recompute_vpin()
composite = composite_signal(
obi=self._obi,
trade_imbalance=self.trade_imbalance(),
vpin=self._current_vpin,
funding_regime=funding_regime,
spread_bps=self._spread_bps,
)
return {
"mid": round(self._mid, 2),
"microprice": round(self._microprice, 2),
"obi": round(self._obi, 4),
"obi_ema": self.ema_obi(),
"obi_std": round(self.obi_volatility(), 4),
"vpin": round(self._current_vpin, 4),
"spread_bps": round(self._spread_bps, 2),
"depth_bid": round(self._depth_bid, 6),
"depth_ask": round(self._depth_ask, 6),
"trade_imbalance": round(self.trade_imbalance(100), 4),
"hft_regime": self.hft_regime(),
"signal": composite["signal"],
"confidence": composite["confidence"],
"breakdown": composite["breakdown"],
"trade_count": self._trade_count,
}
# ── Getters ───────────────────────────────────────────────
@property
def mid(self) -> float:
return self._mid
@property
def obi(self) -> float:
return self._obi
@property
def vpin(self) -> float:
return self._current_vpin
@property
def spread_bps(self) -> float:
return self._spread_bps
+370
View File
@@ -0,0 +1,370 @@
"""
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 live.treasury import Treasury
from live.filters.toxicity import ToxicityFilter
from live.makers.hl_btc_eth import HlMakerPool
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
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)
# 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 — mark-based)
if self._mode == "paper":
for coin in self._coins:
q = quotes.get(coin)
if q:
pipeline = self._pipelines[coin]
self._simulate_paper_fills(coin, q, pipeline)
# 8. Update funding monitor
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):
"""Naive paper fill: if our bid > mid or ask < mid after some random threshold,
simulate a fill. In production this comes from exchange WebSocket."""
import random
mid = pipeline.mid
if mid <= 0:
return
if random.random() < 0.05:
side = "bid" if random.random() < 0.5 else "ask"
size = getattr(quote, f"{side}_size", 0.001)
px = getattr(quote, side, 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)
# ── 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(),
"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())
+144
View File
@@ -0,0 +1,144 @@
"""
Integration tests for the analytics pipeline and production node.
"""
from live.integrator import AnalyticsPipeline
class TestAnalyticsPipeline:
def test_empty_pipeline(self):
p = AnalyticsPipeline()
result = p.emit()
assert result["signal"] == "neutral"
assert result["mid"] == 0
def test_book_update_sets_mid(self):
p = AnalyticsPipeline()
p.update_book({50000.0: 1.0, 49999.0: 2.0}, {50002.0: 1.0, 50003.0: 0.5})
assert p.mid == 50001.0
assert p.obi != 0
assert p.spread_bps > 0
def test_trade_updates_vpin(self):
p = AnalyticsPipeline()
p.update_book({50000.0: 5.0, 49999.0: 5.0}, {50001.0: 5.0, 50002.0: 5.0})
for _ in range(100):
p.update_trade(50001.0, 0.01, 50000.5) # buys
for _ in range(50):
p.update_trade(50000.0, 0.005, 50000.5) # sells
assert p.vpin >= 0
assert p.trade_imbalance() > 0 # more buys
def test_emi_of(self):
p = AnalyticsPipeline()
p.update_book({50000.0: 1.0, 49999.0: 2.0}, {50002.0: 1.0, 50003.0: 0.5})
assert p.ema_obi(alpha=0.5) != 0
def test_full_emit(self):
p = AnalyticsPipeline()
bids = {100.0: 1.0, 99.0: 2.0}
asks = {102.0: 1.0, 103.0: 3.0}
p.update_book(bids, asks)
for _ in range(20):
p.update_trade(101.0, 0.1, 101.0)
result = p.emit(funding_regime="neutral")
assert "signal" in result
assert "confidence" in result
assert "vpin" in result
assert "obi" in result
assert "hft_regime" in result
assert "breakdown" in result
def test_hft_regime_detect(self):
p = AnalyticsPipeline()
bids = {100.0: 1.0}
asks = {102.0: 1.0}
p.update_book(bids, asks)
regime = p.hft_regime()
assert regime in ("trending", "ranging", "toxic", "quiet")
def test_pipeline_per_coin_isolation(self):
btc = AnalyticsPipeline()
eth = AnalyticsPipeline()
btc.update_book({50000.0: 1.0}, {50002.0: 1.0})
eth.update_book({3000.0: 1.0}, {3002.0: 1.0})
assert btc.mid > 40000
assert eth.mid < 10000
def test_trade_count_tracking(self):
p = AnalyticsPipeline()
p.update_book({100.0: 1.0}, {102.0: 1.0})
for _ in range(5):
p.update_trade(101.0, 0.1)
assert p.emit()["trade_count"] == 5
class TestNodeV2Smoke:
def test_node_creation(self):
from live.node_v2 import ProductionNode
node = ProductionNode(
coins=["BTC"],
testnet=True,
mode="paper",
max_position_per_coin=0.001,
base_quote_size=0.0001,
)
assert node is not None
def test_node_start_stop(self):
import asyncio
from live.node_v2 import ProductionNode
async def _test():
node = ProductionNode(
coins=["BTC"],
testnet=True,
mode="paper",
tick_interval_sec=0.1,
max_position_per_coin=0.001,
base_quote_size=0.0001,
)
await node.start()
for _ in range(3):
await node._tick_cycle()
await node.stop()
assert node._tick > 0
asyncio.run(_test())
def test_metrics_written(self):
import asyncio, json, tempfile, os, time
from live.node_v2 import ProductionNode
f = tempfile.NamedTemporaryFile(delete=False, suffix=".json")
f.close()
async def _test():
node = ProductionNode(
coins=["BTC"],
testnet=True,
mode="paper",
tick_interval_sec=0.1,
max_position_per_coin=0.001,
base_quote_size=0.0001,
metrics_file=f.name,
)
await node.start()
for _ in range(2):
await node._tick_cycle()
await node.stop()
assert os.path.exists(f.name)
data = json.load(open(f.name))
assert "treasury" in data
assert "analytics" in data
assert "equity_history" in data
assert "maker" in data
asyncio.run(_test())
os.unlink(f.name)
class TestCLISmoke:
def test_cli_import(self):
import cli
assert cli.main is not None