""" Flow toxicity and adverse selection analytics. VPIN (Volume-synchronized Probability of Informed Trading), fill toxicity metrics, and adverse selection indicators based on order book and trade data. """ from __future__ import annotations import math from collections import deque import numpy as np # ── VPIN ───────────────────────────────────────────────────── def compute_vpin( buy_volume: list[float], sell_volume: list[float], volume_bucket_size: float | None = None, n_buckets: int = 50, ) -> dict: """Volume-synchronized Probability of Informed Trading. Args: buy_volume: volume classified as buyer-initiated per bar/period sell_volume: volume classified as seller-initiated per bar/period volume_bucket_size: target volume per bucket (auto if None) n_buckets: number of buckets for rolling VPIN Returns dict with vpin values and summary. """ if not buy_volume or len(buy_volume) != len(sell_volume): return {"vpin_value": 0, "n_buckets": 0, "bucket_size": 0} if volume_bucket_size is None: total_vol = sum(buy_volume) + sum(sell_volume) volume_bucket_size = total_vol / max(len(buy_volume), 1) * 5 buy_sell = [(b, s) for b, s in zip(buy_volume, sell_volume)] buckets: list[dict] = [] current_buy = 0.0 current_sell = 0.0 for b, s in buy_sell: current_buy += b current_sell += s if current_buy + current_sell >= volume_bucket_size: total = current_buy + current_sell buckets.append({ "buy": current_buy, "sell": current_sell, "total": total, "imbalance": abs(current_buy - current_sell), }) # Carry over excess excess = total - volume_bucket_size current_buy = excess * (current_buy / total) if total > 0 else 0 current_sell = excess * (current_sell / total) if total > 0 else 0 vpin_value = 0.0 if len(buckets) >= n_buckets: recent = buckets[-n_buckets:] total_imb = sum(b["imbalance"] for b in recent) total_vol = sum(b["total"] for b in recent) vpin_value = total_imb / total_vol if total_vol > 0 else 0.0 return { "vpin_value": round(vpin_value, 4), "n_buckets": len(buckets), "bucket_size": round(volume_bucket_size, 2), } def compute_vpin_time_series( buy_volume: list[float], sell_volume: list[float], volume_bucket_size: float | None = None, n_buckets: int = 50, ) -> dict: """Compute rolling VPIN time series.""" if not buy_volume or len(buy_volume) != len(sell_volume): return {"vpin_values": [], "mean": 0, "std": 0, "max": 0, "threshold_alarm": 0} if volume_bucket_size is None: total_vol = sum(buy_volume) + sum(sell_volume) volume_bucket_size = total_vol / max(len(buy_volume), 1) * 5 buy_sell = [(b, s) for b, s in zip(buy_volume, sell_volume)] buckets: list[float] = [] current_buy = 0.0 current_sell = 0.0 vpin_series = [] for b, s in buy_sell: current_buy += b current_sell += s if current_buy + current_sell >= volume_bucket_size: total = current_buy + current_sell imb = abs(current_buy - current_sell) buckets.append(imb / total if total > 0 else 0.5) excess = total - volume_bucket_size ratio = current_buy / total if total > 0 else 0.5 current_buy = ratio * excess current_sell = (1 - ratio) * excess if len(buckets) >= n_buckets: vpin_series.append(sum(buckets[-n_buckets:]) / n_buckets) else: vpin_series.append(sum(buckets) / len(buckets)) a = np.array(vpin_series) if vpin_series else np.array([0.0]) return { "vpin_values": [round(v, 4) for v in vpin_series], "mean": round(float(np.mean(a)), 4), "std": round(float(np.std(a)), 4), "max": round(float(np.max(a)), 4), "threshold_alarm": round(float(np.mean(a) + 2 * np.std(a)), 4), } # ── Fill toxicity ──────────────────────────────────────────── def fill_toxicity( trade_prices: list[float], mids: list[float], trade_sides: list[str], horizon_ticks: int = 10, ) -> dict: """Compute toxicity per trade: did price move against you after fill? For each buy trade: toxicity = (mid before - mid after) / mid For each sell: toxicity = (mid after - mid before) / mid Positive toxicity = adverse price movement post-trade. """ if len(trade_prices) < horizon_ticks + 1: return {"buy_toxicity_mean": 0, "sell_toxicity_mean": 0, "overall": 0} buy_tox = [] sell_tox = [] n = len(trade_prices) for i in range(n - horizon_ticks): side = trade_sides[i] if i < len(trade_sides) else "unknown" mid_before = mids[min(i + 1, n - 1)] mid_after = mids[min(i + horizon_ticks, n - 1)] if mid_before <= 0 or mid_after <= 0: continue change = (mid_before - mid_after) / mid_before if side == "buy": buy_tox.append(change) elif side == "sell": sell_tox.append(-change) return { "buy_toxicity_mean_bps": round(float(np.mean(buy_tox)) * 10000, 2) if buy_tox else 0, "sell_toxicity_mean_bps": round(float(np.mean(sell_tox)) * 10000, 2) if sell_tox else 0, "buy_count": len(buy_tox), "sell_count": len(sell_tox), "overall_bps": round(float(np.mean(buy_tox + sell_tox)) * 10000, 2) if buy_tox or sell_tox else 0, } # ── Adverse selection ─────────────────────────────────────── def adverse_selection_ratio( mid_after_trades: list[float], mid_before_trades: list[float], trade_sides: list[str], ) -> dict: """Adverse selection ratio per side (mid after / mid before - 1). Higher values = more adverse selection (price moves against you). """ buys = [] sells = [] for i, side in enumerate(trade_sides): if i >= len(mid_after_trades) or i >= len(mid_before_trades): break before = mid_before_trades[i] after = mid_after_trades[i] if before <= 0: continue sel = (after - before) / before if side == "buy": buys.append(-sel) elif side == "sell": sells.append(sel) ba = np.array(buys) if buys else np.array([0.0]) sa = np.array(sells) if sells else np.array([0.0]) return { "buy_adverse_bps": round(float(np.mean(ba)) * 10000, 2), "sell_adverse_bps": round(float(np.mean(sa)) * 10000, 2), "buy_win_pct": round(np.sum(ba <= 0) / len(ba), 4) if len(ba) > 0 else 0, "sell_win_pct": round(np.sum(sa <= 0) / len(sa), 4) if len(sa) > 0 else 0, } # ── Liquidation clustering ────────────────────────────────── def liquidation_clustering( liquidation_times_ms: list[int], window_sec: int = 300, ) -> dict: """Detect liquidation clusters — unusual concentration of liquidations. Returns cluster periods and intensity. """ if len(liquidation_times_ms) < 2: return {"clusters": [], "mean_interval_s": 0, "clustered_pct": 0} intervals = [ (liquidation_times_ms[i + 1] - liquidation_times_ms[i]) / 1000 for i in range(len(liquidation_times_ms) - 1) ] mean_interval = float(np.mean(intervals)) std_interval = float(np.std(intervals)) clusters = [] cluster_start = None for i, interval in enumerate(intervals): if interval < mean_interval * 0.3: # Tight clustering threshold if cluster_start is None: cluster_start = liquidation_times_ms[i] else: if cluster_start is not None and liquidation_times_ms[i] - cluster_start < window_sec * 1000: clusters.append({ "start_ms": cluster_start, "end_ms": liquidation_times_ms[i], "count": i - liquidation_times_ms.index(cluster_start) + 1 if cluster_start in liquidation_times_ms else 0, }) cluster_start = None clustered_count = sum(c.get("count", 0) for c in clusters) return { "clusters": clusters, "n_clusters": len(clusters), "mean_interval_s": round(mean_interval, 2), "clustered_pct": round(clustered_count / len(liquidation_times_ms), 4) if liquidation_times_ms else 0, }