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