e58c5951b7
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.
215 lines
7.0 KiB
Python
215 lines
7.0 KiB
Python
"""
|
|
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
|