From 5304534e380c973795ca4335d62f2cfa11c96927 Mon Sep 17 00:00:00 2001 From: ramseshk <45832522+ramseshk@users.noreply.github.com> Date: Fri, 7 Aug 2026 17:52:20 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20advanced=20microstructure=20modules=20?= =?UTF-8?q?=E2=80=94=20HLP,=20Hawkes,=20Whipsaw,=20Term=20Structure,=20Liq?= =?UTF-8?q?=20Waterfall,=20Spoof=20Detector?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 6 new modules with 46 new tests (230 total): #21 HLP Vault Monitor (live/monitors/hlp_vault.py): Tracks Hyperliquid's native protocol market maker at address 0xfefefe... Queries clearinghouseState + metaAndAssetCtxs. - Delta exposure per asset (notional + PnL) - Overextension detection (notional exceeds M threshold) - Rebalancing signals: fade_short when HLP too short, fade_long when HLP too long (front-run forced rebalancing) - Toxicity score: HLP losing money = absorbing informed flow - Historical delta tracking #29 Hawkes Processes (microstructure/hawkes.py): Multivariate Hawkes calibrator for limit order book dynamics. - MLE calibration via SGD gradient descent on log-likelihood - Branching ratio enforcement (alpha/beta < 0.99 for stationarity) - Intensity computation λ_i(t) with cross-excitation - Activity forecasting (expected event count in horizon) - Synthetic event generator (Ogata thinning) - Pure functions: hawkes_intensity, hawkes_log_likelihood, generate_hawkes_events #23 Funding Whipsaw Trader (live/strategies/funding_whipsaw.py): Premium index decay trading in final 60s of funding epoch. - Detects deterministic convergence of premium→0 at settlement - Time-scaled position sizing (larger closer to settlement) - Auto-close after funding epoch completes - Confidence scoring based on premium magnitude #32 Term Structure Monitor (live/monitors/term_structure.py): Perp/quarterly/bi-quarterly futures basis curve trading. - Quarterly-perp basis with z-score anomaly detection - BiQ-quarterly curve steepness monitoring - Fair quarterly price via interest rate parity + funding carry - Calendar spread signals: buy_basis, sell_basis, curve_steepener, curve_flattener #24 Liquidation Waterfall (live/monitors/liq_waterfall.py): Cross-margin liquidation order prediction. - Margin ratio tracking (equity / maintenance margin) - Danger/critical level classification - Asset liquidation priority: maintenance / book_liquidity ratio (least liquid asset relative to margin = dumped first) - Strategy output: widen_spreads on target, tighten on rest #31 Spoof Detector (microstructure/spoof_detector.py): Adversarial ML-style spoofing pattern recognition. - Rule 1: Large order far from mid, cancelled immediately - Rule 2: Cancel right before trade approaches price level - Rule 3: Oversized order with no fill within short lifetime - Spoof probability (rolling window ratio) - Cancel-to-fill ratio monitoring --- live/monitors/hlp_vault.py | 269 ++++++++++++++++++++++++ live/monitors/liq_waterfall.py | 168 +++++++++++++++ live/monitors/term_structure.py | 197 ++++++++++++++++++ live/strategies/__init__.py | 0 live/strategies/funding_whipsaw.py | 195 ++++++++++++++++++ microstructure/hawkes.py | 320 +++++++++++++++++++++++++++++ microstructure/spoof_detector.py | 180 ++++++++++++++++ tests/test_advanced_monitors.py | 170 +++++++++++++++ tests/test_funding_whipsaw.py | 91 ++++++++ tests/test_hawkes.py | 149 ++++++++++++++ tests/test_hlp_vault.py | 153 ++++++++++++++ 11 files changed, 1892 insertions(+) create mode 100644 live/monitors/hlp_vault.py create mode 100644 live/monitors/liq_waterfall.py create mode 100644 live/monitors/term_structure.py create mode 100644 live/strategies/__init__.py create mode 100644 live/strategies/funding_whipsaw.py create mode 100644 microstructure/hawkes.py create mode 100644 microstructure/spoof_detector.py create mode 100644 tests/test_advanced_monitors.py create mode 100644 tests/test_funding_whipsaw.py create mode 100644 tests/test_hawkes.py create mode 100644 tests/test_hlp_vault.py diff --git a/live/monitors/hlp_vault.py b/live/monitors/hlp_vault.py new file mode 100644 index 0000000..5d1023e --- /dev/null +++ b/live/monitors/hlp_vault.py @@ -0,0 +1,269 @@ +""" +HLP (Hyperliquidity Provider) Vault monitoring. + +Tracks Hyperliquid's native protocol-level market-making vault at +address 0xfefefefefefefefefefefefefefefefefefefefe. + +HLP acts as counterparty to all user trades. When it absorbs toxic flow +or becomes directionally overextended, it must rebalance — creating +predictable market impact that can be traded. + +Signals: + - fade_short: HLP is too short → expect buying rebalance → go long + - fade_long: HLP is too long → expect selling rebalance → go short + - neutral: HLP delta is balanced, safe to provide liquidity alongside +""" + +from __future__ import annotations + +import logging +import time +from collections import deque +from typing import Optional + +import requests + +logger = logging.getLogger(__name__) + +HLP_ADDRESS = "0xfefefefefefefefefefefefefefefefefefefefe" +TESTNET_API = "https://api.hyperliquid-testnet.xyz/info" +MAINNET_API = "https://api.hyperliquid.xyz/info" + + +class HlpVaultMonitor: + """Monitor HLP vault state — delta, PnL, rebalancing pressure, toxicity.""" + + def __init__( + self, + testnet: bool = True, + overextended_threshold: float = 5.0, # notional in $M before overextended + history_window: int = 1000, + ): + self._api_url = TESTNET_API if testnet else MAINNET_API + self._overextended_threshold = overextended_threshold * 1_000_000 # convert to USD + self._testnet = testnet + + # Per-coin state + self._positions: dict[str, dict] = {} # coin → {side, szi, entry_px, upnl} + self._mark_prices: dict[str, float] = {} # coin → mark price + self._oracle_prices: dict[str, float] = {} # coin → oracle price + self._funding_rates: dict[str, float] = {} # coin → funding rate + self._open_interest: dict[str, float] = {} # coin → OI + + # History + self._delta_history: dict[str, deque] = {} # coin → deque of (time, notional) + self._pnl_history: list[dict] = [] + self._history_window = history_window + + self._last_update: float = 0 + self._update_count: int = 0 + + # ── Core update ────────────────────────────────────────── + + def _api_post(self, payload: dict) -> dict: + """Make a POST request to HL info API (testable via mock).""" + resp = requests.post(self._api_url, json=payload, timeout=10) + resp.raise_for_status() + return resp.json() + + def update(self): + """Fetch latest HLP state from Hyperliquid API. + + Makes two calls: + 1. metaAndAssetCtxs → mark prices, funding, OI + 2. clearinghouseState(HLP_ADDRESS) → positions, PnL + """ + now = time.time() + + # Fetch market data + try: + meta_and_ctx = self._api_post({"type": "metaAndAssetCtxs"}) + if isinstance(meta_and_ctx, list) and len(meta_and_ctx) >= 2: + universe = meta_and_ctx[0].get("universe", []) + ctxs = meta_and_ctx[1] + for i, asset in enumerate(universe): + name = asset.get("name", "") + if name and i < len(ctxs): + self._mark_prices[name] = float(ctxs[i].get("markPx", 0)) + self._oracle_prices[name] = float(ctxs[i].get("oraclePx", 0)) + self._funding_rates[name] = float(ctxs[i].get("funding", 0)) + self._open_interest[name] = float(ctxs[i].get("openInterest", 0)) + except Exception as e: + logger.warning("HLP meta fetch error: %s", e) + + # Fetch HLP positions + try: + ch_state = self._api_post({ + "type": "clearinghouseState", + "user": HLP_ADDRESS, + }) + self._parse_positions(ch_state, now) + except Exception as e: + logger.warning("HLP clearinghouse fetch error: %s", e) + + self._last_update = now + self._update_count += 1 + + def _parse_positions(self, data: dict, now: float): + """Parse clearinghouseState response into per-coin positions.""" + asset_positions = data.get("assetPositions", []) + new_positions: dict[str, dict] = {} + + for ap in asset_positions: + pos = ap.get("position", {}) + if not pos: + continue + coin = pos.get("coin", "") + if not coin: + continue + szi = float(pos.get("szi", 0)) + entry_px = float(pos.get("entryPx", 0)) + upnl = float(pos.get("unrealizedPnl", 0)) + side = pos.get("side", "") # "A" = short, "B" = long + + # HLP short is side="A", long is side="B" + signed_szi = -szi if side == "A" else szi + + new_positions[coin] = { + "side": side, + "szi": szi, + "signed_szi": signed_szi, + "entry_px": entry_px, + "unrealized_pnl": upnl, + "mark_px": self._mark_prices.get(coin, entry_px), + "notional_usd": abs(szi) * self._mark_prices.get(coin, entry_px), + } + + # Update delta history + if coin not in self._delta_history: + self._delta_history[coin] = deque(maxlen=self._history_window) + self._delta_history[coin].append({ + "t": now, + "signed_szi": signed_szi, + "notional_usd": new_positions[coin]["notional_usd"], + "upnl": upnl, + }) + + self._positions = new_positions + self._pnl_history.append({ + "t": now, + "total_upnl": sum(p["unrealized_pnl"] for p in new_positions.values()), + "asset_count": len(new_positions), + }) + if len(self._pnl_history) > self._history_window: + self._pnl_history = self._pnl_history[-self._history_window:] + + # ── Queries ────────────────────────────────────────────── + + def position(self, coin: str) -> float: + """Signed position for a coin (positive = long).""" + pos = self._positions.get(coin.upper(), {}) + return pos.get("signed_szi", 0.0) + + def delta_exposure(self) -> dict: + """Delta exposure per coin with notional and PnL.""" + result = {} + for coin, pos in self._positions.items(): + result[coin] = { + "signed_size": pos["signed_szi"], + "notional_usd": round(pos["notional_usd"], 2), + "unrealized_pnl": round(pos["unrealized_pnl"], 2), + "side": "long" if pos["side"] == "B" else "short", + } + return result + + def is_overextended(self, coin: str) -> bool: + """Check if HLP delta on this coin exceeds the threshold.""" + pos = self._positions.get(coin.upper(), {}) + notional = pos.get("notional_usd", 0) + return notional > self._overextended_threshold + + def overextended_assets(self) -> list[str]: + """List of assets where HLP is overextended.""" + return [c for c in self._positions if self.is_overextended(c)] + + def rebalancing_signal(self, coin: str) -> dict: + """Generate a trading signal based on HLP rebalancing pressure. + + Returns: + signal: "fade_short" | "fade_long" | "neutral" + direction: -1 (short) | 0 | 1 (long) for the TRADE direction + overextended: whether HLP is over capacity + """ + pos = self._positions.get(coin.upper(), {}) + if not pos: + return {"signal": "neutral", "direction": 0, "overextended": False, "reason": "no_position"} + + signed = pos["signed_szi"] + is_over = self.is_overextended(coin) + + if not is_over: + return {"signal": "neutral", "direction": 0, "overextended": False, "reason": "balanced"} + + # HLP is short (side=A, signed_szi negative) → it will need to buy to rebalance + # We should fade the short (go long) + if signed < -0.001: + return { + "signal": "fade_short", + "direction": 1, + "overextended": True, + "reason": f"HLP short {abs(signed):.2f} units, expect buying rebalance", + "notional_usd": round(pos["notional_usd"], 2), + } + # HLP is long (side=B, signed_szi positive) → it will need to sell to rebalance + # We should fade the long (go short) + elif signed > 0.001: + return { + "signal": "fade_long", + "direction": -1, + "overextended": True, + "reason": f"HLP long {signed:.2f} units, expect selling rebalance", + "notional_usd": round(pos["notional_usd"], 2), + } + else: + return {"signal": "neutral", "direction": 0, "overextended": False, "reason": "flat"} + + def toxicity_score(self) -> float: + """Estimate how much toxic flow HLP is absorbing. + + Higher score = HLP is losing money = informed traders are beating it. + Range: 0 (healthy) to 1 (toxic). + + Uses: unrealized PnL / total notional as a proxy. + """ + total_notional = sum(p["notional_usd"] for p in self._positions.values()) + total_upnl = sum(p["unrealized_pnl"] for p in self._positions.values()) + + if total_notional <= 0: + return 0.0 + + # Negative PnL → toxic score > 0 + # Positive PnL → toxic score 0 (healthy) + loss_ratio = max(0.0, -total_upnl / total_notional) + toxicity = min(1.0, loss_ratio * 10) # scale: 10% loss = 1.0 toxicity + return round(toxicity, 4) + + def delta_history(self, coin: str) -> list[dict]: + """Historical delta trace for a coin.""" + return list(self._delta_history.get(coin.upper(), [])) + + # ── Summary ────────────────────────────────────────────── + + def summary(self) -> dict: + """One-shot summary of HLP state for dashboard/monitoring.""" + assets = list(self._positions.keys()) + total_delta = sum(p["notional_usd"] for p in self._positions.values()) + overextended = self.overextended_assets() + signals = {coin: self.rebalancing_signal(coin) for coin in assets} + + return { + "assets_tracked": len(assets), + "total_delta_usd": round(total_delta, 2), + "total_delta_m": round(total_delta / 1_000_000, 2), + "toxicity_score": self.toxicity_score(), + "overextended_assets": overextended, + "signals": signals, + "positions": self.delta_exposure(), + "last_update": self._last_update, + "update_count": self._update_count, + } diff --git a/live/monitors/liq_waterfall.py b/live/monitors/liq_waterfall.py new file mode 100644 index 0000000..f60079e --- /dev/null +++ b/live/monitors/liq_waterfall.py @@ -0,0 +1,168 @@ +""" +Cross-margin liquidation waterfall prediction. + +When a whale's cross-margin portfolio approaches liquidation, the +Hyperliquid liquidation engine selects which asset to dump first based +on maintenance margin requirements and order book liquidity. + +By monitoring large cross-margin accounts via clearinghouseState, +we can predict which asset gets liquidated first and position accordingly: + - Widen spreads on the predicted liquidation asset + - Tighten spreads on non-liquidation assets + - Pre-position for the post-liquidation bounce +""" + +from __future__ import annotations + +import logging +from collections import deque +from typing import Optional + +logger = logging.getLogger(__name__) + + +class LiquidationWaterfall: + """Predict the order of cross-margin liquidations for large accounts.""" + + def __init__( + self, + danger_margin_ratio: float = 1.2, # margin_ratio < this = danger + critical_margin_ratio: float = 1.05, # margin_ratio < this = imminent + maintenance_margin_pct: float = 0.03, # 3% maintenance + book_depth_window: int = 10, # levels to estimate liquidity + ): + self._danger = danger_margin_ratio + self._critical = critical_margin_ratio + self._mm_pct = maintenance_margin_pct + self._depth_window = book_depth_window + + self._accounts: dict[str, dict] = {} # address → positions, equity + self._book_depths: dict[str, dict] = {} # coin → {bid_depth, ask_depth} + self._history: deque = deque(maxlen=1000) + + # ── Data feed ──────────────────────────────────────────── + + def update_account(self, address: str, positions: list[dict], margin_balance: float): + """Update a tracked account's positions and margin.""" + self._accounts[address] = { + "positions": positions, + "margin_balance": margin_balance, + "margin_ratio": self._compute_margin_ratio(positions, margin_balance), + } + + def update_book_depth(self, coin: str, bid_depth: float, ask_depth: float): + """Update estimated order book depth for a coin.""" + self._book_depths[coin.upper()] = { + "bid_depth": bid_depth, + "ask_depth": ask_depth, + } + + # ── Risk assessment ───────────────────────────────────── + + def _compute_margin_ratio( + self, positions: list[dict], margin_balance: float + ) -> float: + """Compute margin ratio = equity / maintenance_margin.""" + if not positions or margin_balance <= 0: + return float("inf") + + total_mm = 0.0 + for pos in positions: + size = abs(float(pos.get("szi", 0))) + px = float(pos.get("entryPx", pos.get("markPx", 0))) + total_mm += size * px * self._mm_pct + + return margin_balance / total_mm if total_mm > 0 else float("inf") + + def at_risk_accounts(self) -> list[dict]: + """List accounts approaching liquidation.""" + risky = [] + for addr, acct in self._accounts.items(): + ratio = acct["margin_ratio"] + if ratio < self._danger: + level = "critical" if ratio < self._critical else "danger" + risky.append({ + "address": addr[:10] + "...", + "margin_ratio": round(ratio, 3), + "level": level, + "positions": len(acct["positions"]), + }) + return sorted(risky, key=lambda r: r["margin_ratio"]) + + def predict_liquidation_order(self, address: str) -> list[dict]: + """Predict which assets get liquidated first for a given account. + + Returns assets ranked by liquidation priority (first to go = top). + Uses: maintenance_margin_requirement / book_liquidity ratio. + Higher ratio = less liquid relative to margin cost = dumped first. + """ + acct = self._accounts.get(address) + if not acct: + return [] + + positions = acct["positions"] + ranked = [] + + for pos in positions: + coin = pos.get("coin", "").upper() + size = abs(float(pos.get("szi", 0))) + px = float(pos.get("entryPx", pos.get("markPx", 0))) + + maintenance = size * px * self._mm_pct + + # Liquidity: how much the book can absorb before significant slippage + book = self._book_depths.get(coin, {}) + side = "buy" if float(pos.get("szi", 0)) < 0 else "sell" + depth = book.get("bid_depth" if side == "buy" else "ask_depth", size * px * 0.1) + + # Liquidation priority score: higher = dumped first + liquidity_ratio = maintenance / max(depth, 1e-8) + ranked.append({ + "coin": coin, + "size": size, + "notional_usd": round(size * px, 2), + "maintenance_usd": round(maintenance, 2), + "liquidity_ratio": round(liquidity_ratio, 4), + "predicted_first": False, # set below + }) + + # Sort by liquidity ratio (highest = least liquid = dumped first) + ranked.sort(key=lambda r: r["liquidity_ratio"], reverse=True) + if ranked: + ranked[0]["predicted_first"] = True + + return ranked + + def signal(self, address: str) -> dict: + """Generate trading signal based on predicted liquidation waterfall. + + If an account is in danger and we can predict the liquidation order: + - Asset predicted to be dumped first → widen spreads, go short + - Other assets in the portfolio → tighten spreads (safer to quote) + """ + acct = self._accounts.get(address) + if not acct or acct["margin_ratio"] > self._danger: + return {"action": "none", "reason": "account_safe"} + + order = self.predict_liquidation_order(address) + if not order: + return {"action": "none", "reason": "no_positions"} + + first_asset = order[0] + return { + "action": "position", + "reason": f"liquidation_imminent_{acct['margin_ratio']:.2f}", + "margin_ratio": round(acct["margin_ratio"], 3), + "liquidation_target": first_asset["coin"], + "strategy": { + first_asset["coin"]: "widen_spreads_2x", + **{r["coin"]: "tighten_spreads" for r in order[1:]}, + }, + "predicted_order": [r["coin"] for r in order], + } + + def summary(self) -> dict: + return { + "at_risk_accounts": self.at_risk_accounts(), + "tracked_accounts": len(self._accounts), + } diff --git a/live/monitors/term_structure.py b/live/monitors/term_structure.py new file mode 100644 index 0000000..7bca781 --- /dev/null +++ b/live/monitors/term_structure.py @@ -0,0 +1,197 @@ +""" +Term structure monitor — perp vs quarterly vs bi-quarterly futures basis. + +Hyperliquid offers perpetual (funding-based), quarterly, and bi-quarterly +futures contracts. The basis curve (perp→quarterly→bi-quarterly) contains +information about market expectations and can be traded. + +Anomalies: + - Perp funding deeply negative but quarterly basis remains steep → + go long perp (collect funding), short quarterly (lock basis) + - Quarterly futures converging to perp at expiration → + calendar spread mean-reversion + - Bi-quarterly premium over quarterly deviating from fair value → + curve steepener/flattener trades +""" + +from __future__ import annotations + +from collections import deque +from typing import Optional + + +class TermStructureMonitor: + """Monitor perp/futures term structure for arbitrage opportunities. + + Tracks: + - Perp funding rate and mark price + - Quarterly futures price + - Bi-quarterly futures price (if available) + - Basis spreads: quarterly-perp, biq-quarterly, biq-perp + - Calendar spread mean-reversion + """ + + def __init__( + self, + funding_window: int = 1440, # 24h of 1-min samples + basis_window: int = 100, + ): + self._funding_window = funding_window + self._basis_window = basis_window + + self._perp_prices: dict[str, deque[float]] = {} + self._quarterly_prices: dict[str, deque[float]] = {} + self._biq_prices: dict[str, deque[float]] = {} + self._funding_rates: dict[str, deque[float]] = {} + + self._funding_epoch_seconds: int = 8 * 3600 + self._quarterly_expiry_days: int = 90 + self._biq_expiry_days: int = 180 + + # ── Data feed ──────────────────────────────────────────── + + def update_perp(self, coin: str, price: float, funding_rate: float): + c = coin.upper() + self._perp_prices.setdefault(c, deque(maxlen=self._basis_window)).append(price) + self._funding_rates.setdefault(c, deque(maxlen=self._funding_window)).append(funding_rate) + + def update_quarterly(self, coin: str, price: float): + self._quarterly_prices.setdefault(coin.upper(), deque(maxlen=self._basis_window)).append(price) + + def update_biq(self, coin: str, price: float): + self._biq_prices.setdefault(coin.upper(), deque(maxlen=self._basis_window)).append(price) + + # ── Basis computation ──────────────────────────────────── + + def quarterly_perp_basis(self, coin: str) -> dict | None: + """Basis between quarterly future and perpetual.""" + q = list(self._quarterly_prices.get(coin.upper(), [])) + p = list(self._perp_prices.get(coin.upper(), [])) + min_len = min(len(q), len(p)) + if min_len < 2: + return None + + qq = q[-min_len:] + pp = p[-min_len:] + basis_bps = [(qq[i] - pp[i]) / pp[i] * 10000 for i in range(min_len) if pp[i] > 0] + + if not basis_bps: + return None + + return { + "current_bps": round(basis_bps[-1], 2), + "mean_bps": round(sum(basis_bps) / len(basis_bps), 2), + "std_bps": round(_std(basis_bps), 2), + "z_score": round((basis_bps[-1] - sum(basis_bps) / len(basis_bps)) / max(_std(basis_bps), 0.01), 2), + "n_samples": len(basis_bps), + } + + def biq_quarterly_basis(self, coin: str) -> dict | None: + """Basis between bi-quarterly and quarterly futures (curve steepness).""" + bq = list(self._biq_prices.get(coin.upper(), [])) + q = list(self._quarterly_prices.get(coin.upper(), [])) + min_len = min(len(bq), len(q)) + if min_len < 2: + return None + + bb = bq[-min_len:] + qq = q[-min_len:] + basis_bps = [(bb[i] - qq[i]) / qq[i] * 10000 for i in range(min_len) if qq[i] > 0] + + if not basis_bps: + return None + + return { + "current_bps": round(basis_bps[-1], 2), + "mean_bps": round(sum(basis_bps) / len(basis_bps), 2), + "std_bps": round(_std(basis_bps), 2), + "z_score": round((basis_bps[-1] - sum(basis_bps) / len(basis_bps)) / max(_std(basis_bps), 0.01), 2), + "n_samples": len(basis_bps), + } + + def fair_quarterly_price(self, coin: str, risk_free_annual: float = 0.05) -> dict | None: + """Compute fair quarterly price from perp via interest rate parity.""" + p = list(self._perp_prices.get(coin.upper(), [])) + if not p: + return None + + perp_px = p[-1] + days = self._quarterly_expiry_days + funding = list(self._funding_rates.get(coin.upper(), [])) + avg_funding = sum(funding[-100:]) / max(len(funding[-100:]), 1) if funding else 0 + + # Fair quarterly = perp * (1 + (r + avg_funding) * days/365) + carry_rate = risk_free_annual + avg_funding * 3 * 365 # annualize 8h funding + fair_px = perp_px * (1 + carry_rate * days / 365) + + return { + "perp_price": perp_px, + "fair_quarterly": round(fair_px, 2), + "carry_rate_annual_pct": round(carry_rate * 100, 2), + "days_to_expiry": days, + } + + def signal(self, coin: str) -> dict: + """Generate term-structure trading signal. + + Returns: + signal: "buy_basis", "sell_basis", "curve_steepener", "curve_flattener", "none" + """ + c = coin.upper() + qp = self.quarterly_perp_basis(c) + bq = self.biq_quarterly_basis(c) + fair = self.fair_quarterly_price(c) + + signals = [] + + # Check quarterly-perp basis anomalies + if qp and abs(qp["z_score"]) > 2.0: + if qp["z_score"] > 0: + signals.append({ + "signal": "sell_basis", + "reason": f"Quarterly {qp['z_score']:.1f}σ rich vs perp", + "confidence": min(1.0, abs(qp["z_score"]) / 4.0), + }) + else: + signals.append({ + "signal": "buy_basis", + "reason": f"Quarterly {qp['z_score']:.1f}σ cheap vs perp", + "confidence": min(1.0, abs(qp["z_score"]) / 4.0), + }) + + # Check curve steepness + if bq and abs(bq["z_score"]) > 2.0: + if bq["z_score"] > 0: + signals.append({ + "signal": "curve_flattener", + "reason": f"BiQ {bq['z_score']:.1f}σ rich vs quarterly", + "confidence": min(1.0, abs(bq["z_score"]) / 4.0), + }) + else: + signals.append({ + "signal": "curve_steepener", + "reason": f"BiQ {bq['z_score']:.1f}σ cheap vs quarterly", + "confidence": min(1.0, abs(bq["z_score"]) / 4.0), + }) + + result = { + "coin": c, + "signals": signals, + "primary_signal": signals[0]["signal"] if signals else "none", + "quarterly_perp_basis": qp, + "biq_quarterly_basis": bq, + "fair_quarterly": fair, + } + + return result + + def summary(self) -> dict: + return {coin: self.signal(coin) for coin in self._perp_prices} + + +def _std(vals: list[float]) -> float: + """Population standard deviation.""" + 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 diff --git a/live/strategies/__init__.py b/live/strategies/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/live/strategies/funding_whipsaw.py b/live/strategies/funding_whipsaw.py new file mode 100644 index 0000000..6bade25 --- /dev/null +++ b/live/strategies/funding_whipsaw.py @@ -0,0 +1,195 @@ +""" +Funding rate whipsaw trader — premium index decay in final seconds of funding epoch. + +Hyperliquid funding settles every 8 hours (UTC 00:00, 08:00, 16:00). +In the final 60 seconds before settlement, the premium index (Mark - Oracle) +must converge to prevent arbitrage. HFTs trade this convergence deterministically. + +The strategy: + - If premium is positive with <60s until funding → SHORT perp (price will drop) + - If premium is negative with <60s until funding → LONG perp (price will rise) + - Scale position based on premium magnitude and time remaining + - Close position at funding settlement (T+0) +""" +from __future__ import annotations + +import time +from datetime import datetime, timezone +from typing import Optional + + +class FundingWhipsawTrader: + """Trade the deterministic decay of the premium index in the final seconds + before Hyperliquid's funding settlement. + + Funding epochs: 00:00, 08:00, 16:00 UTC every day. + Premium index = (mark_price - oracle_price) / oracle_price + + Usage: + trader = FundingWhipsawTrader() + trader.update(mark_px=64500, oracle_px=64480) + signal = trader.signal() + if signal['action'] != 'none': + # place order: signal['side'], signal['size'], signal['confidence'] + """ + + def __init__( + self, + min_premium_bps: float = 0.5, # minimum premium in bps to trigger + max_size: float = 0.001, # max position size + enter_seconds_before: float = 60.0, # seconds before funding to enter + close_seconds_after: float = 5.0, # seconds after funding to close + ): + self._min_premium_bps = min_premium_bps + self._max_size = max_size + self._enter_seconds = enter_seconds_before + self._close_seconds = close_seconds_after + + self._mark_px: float = 0.0 + self._oracle_px: float = 0.0 + self._premium_bps: float = 0.0 + self._last_update: float = 0.0 + self._in_position: bool = False + self._position_side: str = "" + self._entry_time: float = 0.0 + + def update(self, mark_px: float, oracle_px: float): + """Feed current mark and oracle prices.""" + self._mark_px = mark_px + self._oracle_px = oracle_px + self._last_update = time.time() + + if oracle_px > 0: + self._premium_bps = (mark_px - oracle_px) / oracle_px * 10000 + + def seconds_to_funding(self) -> float: + """Seconds until the next funding settlement (every 8 hours, UTC).""" + now = datetime.now(timezone.utc) + epoch_hours = [0, 8, 16] # Funding at 00:00, 08:00, 16:00 UTC + current_hour = now.hour + + # Find next funding epoch + next_epoch_hour = None + for h in epoch_hours: + if h > current_hour or (h == current_hour and now.minute == 0 and now.second < 5): + next_epoch_hour = h + break + + if next_epoch_hour is None: + # After 16:00, next is 00:00 tomorrow + next_epoch = now.replace(hour=0, minute=0, second=0, microsecond=0) + from datetime import timedelta + next_epoch += timedelta(days=1) + else: + next_epoch = now.replace(hour=next_epoch_hour, minute=0, second=0, microsecond=0) + + delta = (next_epoch - now).total_seconds() + return max(0.0, delta) + + def seconds_since_funding(self) -> float: + """Seconds since the most recent funding settlement.""" + seconds_to = self.seconds_to_funding() + if seconds_to < 3600: # <1 hour to next + return 8 * 3600 - seconds_to + return 8 * 3600 + (3600 - seconds_to % 3600) # Approximate + + def signal(self) -> dict: + """Generate trading signal based on premium and time to funding. + + Returns: + action: "enter_long", "enter_short", "close", "none" + side: "buy" or "sell" (for orders) + size: position size (scaled by time remaining) + confidence: 0-1 confidence in the signal + premium_bps: current premium in bps + seconds_to_funding: time until settlement + """ + secs = self.seconds_to_funding() + premium = self._premium_bps + + # After funding + small delay: close any position + secs_since = self.seconds_since_funding() + if self._in_position and secs_since < self._close_seconds: + self._in_position = False + return { + "action": "close", + "side": "sell" if self._position_side == "buy" else "buy", + "size": self._max_size, + "confidence": 1.0, + "premium_bps": round(premium, 2), + "seconds_to_funding": round(secs, 1), + "reason": "funding_settled", + } + + # Outside entry window: no action + if secs > self._enter_seconds or secs < 0: + return { + "action": "none", + "side": "", + "size": 0.0, + "confidence": 0.0, + "premium_bps": round(premium, 2), + "seconds_to_funding": round(secs, 1), + "reason": "outside_entry_window", + } + + # Already in position + if self._in_position: + return { + "action": "hold", + "side": self._position_side, + "size": self._max_size, + "confidence": 0.8, + "premium_bps": round(premium, 2), + "seconds_to_funding": round(secs, 1), + "reason": "holding", + } + + # Check premium threshold + if abs(premium) < self._min_premium_bps: + return { + "action": "none", + "side": "", + "size": 0.0, + "confidence": 0.0, + "premium_bps": round(premium, 2), + "seconds_to_funding": round(secs, 1), + "reason": "premium_too_small", + } + + # Scale size by time remaining (more remaining = more uncertainty = smaller size) + time_factor = max(0.3, secs / self._enter_seconds) + scaled_size = self._max_size * (1.0 - time_factor * 0.5) + + if premium > 0: + # Premium positive → perp is expensive → short it + side = "sell" + action = "enter_short" + confidence = min(1.0, abs(premium) / self._min_premium_bps * 0.3) + else: + # Premium negative → perp is cheap → long it + side = "buy" + action = "enter_long" + confidence = min(1.0, abs(premium) / self._min_premium_bps * 0.3) + + self._in_position = True + self._position_side = side + self._entry_time = time.time() + + return { + "action": action, + "side": side, + "size": round(scaled_size, 8), + "confidence": round(confidence, 3), + "premium_bps": round(premium, 2), + "seconds_to_funding": round(secs, 1), + "reason": f"premium_{premium:.1f}bps_{secs:.0f}s", + } + + @property + def premium_bps(self) -> float: + return self._premium_bps + + @property + def in_position(self) -> bool: + return self._in_position diff --git a/microstructure/hawkes.py b/microstructure/hawkes.py new file mode 100644 index 0000000..883821a --- /dev/null +++ b/microstructure/hawkes.py @@ -0,0 +1,320 @@ +""" +Multivariate Hawkes process calibrator. + +Models self-exciting and cross-exciting point processes for +limit order book events: trades, cancellations, spread widenings. + +A type-j event at time t_k excites the intensity of type-i events: + λ_i(t) = μ_i + Σ_j α_ij * Σ_{t_k < t} exp(-β * (t - t_k)) + +Key applications: + - Queue depletion probability (trade → more trades) + - Cancel cascade detection (cancel → more cancels) + - Spread widening prediction (trade → spread widening) + - Toxicity anticipation (flow → adverse selection) + +Calibration: maximum likelihood via gradient descent on synthetic or +real event data. Braning ratio Σ_j α_ij / β < 1 for stationarity. +""" + +from __future__ import annotations + +import numpy as np + + +class HawkesCalibrator: + """Multivariate Hawkes process with MLE calibration.""" + + def __init__(self, n_dimensions: int = 3, beta: float = 2.0): + self._n_dim = n_dimensions + self._mu = np.ones(n_dimensions) * 0.5 # baseline intensity + self._alpha = np.eye(n_dimensions) * 0.1 # excitation matrix + self._beta = beta # decay rate + + # ── Properties ─────────────────────────────────────────── + + @property + def n_dim(self) -> int: + return self._n_dim + + @property + def mu(self) -> np.ndarray: + return self._mu.copy() + + @mu.setter + def mu(self, value: np.ndarray): + self._mu = np.asarray(value, dtype=float) + + @property + def alpha(self) -> np.ndarray: + return self._alpha.copy() + + @alpha.setter + def alpha(self, value: np.ndarray): + self._alpha = np.asarray(value, dtype=float) + + @property + def beta(self) -> float: + return self._beta + + @beta.setter + def beta(self, value: float): + self._beta = float(value) + + # ── Calibration ────────────────────────────────────────── + + def calibrate( + self, + events: list[tuple[int, float]], + max_time: float, + learning_rate: float = 0.01, + iterations: int = 200, + regularization: float = 0.001, + ): + """Calibrate parameters via stochastic gradient descent on log-likelihood. + + Args: + events: list of (type, timestamp) tuples, must be sorted by time + max_time: total observation window + learning_rate: SGD step size + iterations: number of gradient steps + regularization: L2 penalty on alpha and mu + """ + events_by_type = [[] for _ in range(self._n_dim)] + for ev_type, ev_time in events: + if 0 <= ev_type < self._n_dim: + events_by_type[ev_type].append(ev_time) + + for _ in range(iterations): + grad_mu, grad_alpha = _compute_gradient( + events_by_type, self._mu, self._alpha, self._beta, max_time + ) + # Gradient ascent with L2 regularization + self._mu += learning_rate * (grad_mu - regularization * self._mu) + self._alpha += learning_rate * (grad_alpha - regularization * self._alpha) + # Project to valid range + self._mu = np.maximum(self._mu, 0.001) + self._alpha = np.maximum(self._alpha, 0.0) + # Enforce stationarity: branching ratio < 0.99 + row_sums = self._alpha.sum(axis=1) + for i in range(self._n_dim): + if row_sums[i] > self._beta * 0.99: + self._alpha[i] *= (self._beta * 0.99) / row_sums[i] + + # ── Intensity ──────────────────────────────────────────── + + def intensity( + self, + dim: int, + event_history: list[tuple[int, float]], + current_time: float, + ) -> float: + """Compute intensity λ_i(t) given event history.""" + full = hawkes_intensity( + self._mu, self._alpha, self._beta, + event_history, current_time, dim=int(dim), + ) + return float(full) + + # ── Metrics ────────────────────────────────────────────── + + def branching_ratio(self) -> float: + """Average branching ratio = max_i(Σ_j α_ij / β). Must be < 1.""" + ratios = self._alpha.sum(axis=1) / self._beta + return float(np.max(ratios)) + + def forecast_activity( + self, + events: list[tuple[int, float]], + current_time: float, + horizon: float = 1.0, + n_samples: int = 100, + ) -> np.ndarray: + """Forecast expected event count per type in the next horizon. + + Uses the branching structure: E[N_i] = μ_i * horizon + Σ_j α_ij/β * current_excitation. + """ + expected = np.zeros(self._n_dim) + contribution = np.zeros(self._n_dim) + + # Baseline contribution + expected += self._mu * horizon + + # Excitation from past events + for ev_type, ev_time in events: + if ev_time >= current_time: + continue + decay = np.exp(-self._beta * (current_time - ev_time)) + remaining = (1.0 - np.exp(-self._beta * horizon)) / self._beta + for i in range(self._n_dim): + contribution[i] += self._alpha[i, ev_type] * decay * remaining + + result = expected + contribution + return np.maximum(result, 0.0) + + +# ── Pure functions ────────────────────────────────────────── + +def hawkes_intensity( + mu: np.ndarray, + alpha: np.ndarray, + beta: float, + event_history: list[tuple[int, float]], + current_time: float, + dim: int | None = None, +) -> np.ndarray: + """Vectorized Hawkes intensity. + + λ_i(t) = μ_i + Σ_j α_ij * Σ_{t_k < t} exp(-β(t - t_k)) + + If dim is specified, returns scalar for that dimension. + """ + n_dim = len(mu) + intensity = mu.copy().astype(float) + + for ev_type, ev_time in event_history: + if ev_time >= current_time or ev_type >= n_dim: + continue + decay = np.exp(-beta * (current_time - ev_time)) + intensity += alpha[:, ev_type] * decay + + if dim is not None: + return intensity[dim] + return intensity + + +def hawkes_log_likelihood( + events: list[tuple[int, float]], + mu: np.ndarray, + alpha: np.ndarray, + beta: float, + max_time: float, +) -> float: + """Compute log-likelihood of observing these events under given parameters.""" + n_dim = len(mu) + events_by_type = [[] for _ in range(n_dim)] + for ev_type, ev_time in events: + if 0 <= ev_type < n_dim: + events_by_type[ev_type].append(ev_time) + + # Term 1: sum over events log(λ_i(t_k)) + event_ll = 0.0 + for ev_type, ev_time in events: + lam = mu[ev_type] + for past_type, past_time in events: + if past_time >= ev_time: + break + lam += alpha[ev_type, past_type] * np.exp(-beta * (ev_time - past_time)) + if lam > 0: + event_ll += np.log(lam) + + # Term 2: -∫ λ(t) dt (compensator) + integral = max_time * mu.sum() + + for i in range(n_dim): + for j in range(n_dim): + for t_j in events_by_type[j]: + integral += alpha[i, j] * (1.0 - np.exp(-beta * (max_time - t_j))) / beta + + return event_ll - integral + + +def generate_hawkes_events( + mu: np.ndarray, + alpha: np.ndarray, + beta: float, + max_time: float, + seed: int | None = None, +) -> list[tuple[int, float]]: + """Generate synthetic events from a multivariate Hawkes process (Ogata thinning). + + Returns list of (type, timestamp) sorted by time. + """ + rng = np.random.RandomState(seed) if seed is not None else np.random + n_dim = len(mu) + + # Upper bound for total intensity + lambda_bar = mu.sum() * 1.5 # conservative upper bound + + events: list[tuple[int, float]] = [] + t = 0.0 + + while t < max_time: + # Generate candidate via Poisson with rate lambda_bar (thinning) + t += rng.exponential(1.0 / lambda_bar) if lambda_bar > 0 else max_time + + if t >= max_time: + break + + # Accept with probability λ(t) / lambda_bar + current_intensity = hawkes_intensity(mu, alpha, beta, events, t) + total_intensity = current_intensity.sum() + + if total_intensity / lambda_bar > rng.uniform(0, 1): + # Accept: determine event type proportional to intensity + probs = current_intensity / total_intensity + ev_type = rng.choice(n_dim, p=probs) + events.append((int(ev_type), t)) + + return events + + +# ── Gradient computation ─────────────────────────────────── + +def _compute_gradient( + events_by_type: list[list[float]], + mu: np.ndarray, + alpha: np.ndarray, + beta: float, + max_time: float, +) -> tuple[np.ndarray, np.ndarray]: + """Compute gradient of log-likelihood wrt mu and alpha.""" + n_dim = len(mu) + grad_mu = np.zeros(n_dim) + grad_alpha = np.zeros((n_dim, n_dim)) + + # Build flat event list for gradient computation + all_events = [] + for i, times in enumerate(events_by_type): + for t in times: + all_events.append((i, t)) + all_events.sort(key=lambda x: x[1]) + + # Gradient of log-likelihood + for ev_type, ev_time in all_events: + lam = mu[ev_type] + past_contributions = {} + for pt, ptime in all_events: + if ptime >= ev_time: + break + contrib = alpha[ev_type, pt] * np.exp(-beta * (ev_time - ptime)) + lam += contrib + past_contributions[pt] = contrib + + if lam > 0: + grad_mu[ev_type] += 1.0 / lam + + # Gradient of alpha: Σ_{t_k} α_{ij} * exp(-β (t - t_k)) contribution + for i in range(n_dim): + for j in range(n_dim): + for t_j in events_by_type[j]: + decay_integral = (1.0 - np.exp(-beta * (max_time - t_j))) / beta + grad_alpha[i, j] -= decay_integral + + # Add per-event alpha gradient + for ev_type, ev_time in all_events: + lam = mu[ev_type] + contributions = [] + for pt, ptime in all_events: + if ptime >= ev_time: + break + lam += alpha[ev_type, pt] * np.exp(-beta * (ev_time - ptime)) + + if lam > 0: + for pt, ptime in all_events: + if ptime >= ev_time: + break + contrib = np.exp(-beta * (ev_time - ptime)) / lam + grad_alpha[ev_type, pt] += contrib + + return grad_mu, grad_alpha diff --git a/microstructure/spoof_detector.py b/microstructure/spoof_detector.py new file mode 100644 index 0000000..2cf3f0f --- /dev/null +++ b/microstructure/spoof_detector.py @@ -0,0 +1,180 @@ +""" +Adversarial ML spoof detection — recognize market manipulation patterns +in L3 (order-by-order) data. + +Detects: + 1. Spoofing: large orders placed far from mid, cancelled before execution + 2. Layering: multiple orders at different price levels on one side, + all cancelled simultaneously when price moves + 3. Quote stuffing: rapid order submission and cancellation to slow competitors + 4. Momentum ignition: small aggressive trades followed by large passive orders + +Uses lightweight feature engineering (no deep learning required): + - Order lifetime before cancellation + - Distance from mid price + - Size relative to typical trade size + - Correlation between cancel events and price moves + - Pattern matching on order sequences + +Output feeds into ToxicityFilter for pre-trade gating. +""" + +from __future__ import annotations + +from collections import deque +from typing import Optional + + +class SpoofDetector: + """Detect spoofing patterns in order book event streams. + + Maintains a rolling window of order events (place, cancel, modify) + and classifies each order as legitimate or suspicious. + + Usage: + detector = SpoofDetector() + detector.record_place(order_id, side, price, size, mid, timestamp) + detector.record_cancel(order_id, mid, timestamp) + score = detector.spoof_probability() # 0-1 + if score > 0.5: + # increase toxicity filter, reduce quote sizes + """ + + def __init__( + self, + window_seconds: float = 60.0, + max_orders: int = 1000, + spoof_cancel_threshold: float = 0.5, # % lifetime below mid-distance to flag + size_multiple: float = 3.0, # order size / avg trade size > this = large + price_ticks_threshold: int = 5, # cancel when price moves within N ticks of order + ): + self._window = window_seconds + self._max_orders = max_orders + self._cancel_threshold = spoof_cancel_threshold + self._size_multiple = size_multiple + self._ticks_threshold = price_ticks_threshold + + self._orders: dict[str, dict] = {} # order_id → {side, px, sz, mid_at_place, time} + self._cancel_events: deque = deque(maxlen=max_orders) + self._fill_events: deque = deque(maxlen=max_orders // 2) + self._mid_prices: deque[float] = deque(maxlen=500) + self._trade_sizes: deque[float] = deque(maxlen=500) + + self._spoof_count: int = 0 + self._total_orders: int = 0 + self._total_cancels: int = 0 + + # ── Event recording ────────────────────────────────────── + + def record_place( + self, order_id: str, side: str, price: float, size: float, mid: float, timestamp: float + ): + """Record a new limit order placement.""" + self._orders[order_id] = { + "side": side, + "px": price, + "sz": size, + "mid_at_place": mid, + "time": timestamp, + } + self._total_orders += 1 + self._mid_prices.append(mid) + self._trade_sizes.append(size) + + # Cleanup old orders + if len(self._orders) > self._max_orders: + cutoff = timestamp - self._window + stale = [oid for oid, o in self._orders.items() if o["time"] < cutoff] + for oid in stale: + del self._orders[oid] + + def record_cancel(self, order_id: str, mid: float, timestamp: float): + """Record a cancellation. Returns True if classified as spoof.""" + self._total_cancels += 1 + order = self._orders.pop(order_id, None) + if not order: + self._cancel_events.append({"spoof": False, "time": timestamp}) + return False + + lifetime = timestamp - order["time"] + dist_bps = abs(order["px"] - order["mid_at_place"]) / order["mid_at_place"] * 10000 \ + if order["mid_at_place"] > 0 else 0 + + # Spoof classification rules + is_spoof = False + reasons = [] + + # Rule 1: Large order far from mid, cancelled quickly + avg_size = sum(self._trade_sizes) / max(len(self._trade_sizes), 1) + if order["sz"] > avg_size * self._size_multiple and dist_bps > 20: + if lifetime < self._cancel_threshold * dist_bps: # proportional to distance + is_spoof = True + reasons.append("large_far_quick_cancel") + + # Rule 2: Cancel right before price approaches (within N ticks) + price_moved = abs(mid - order["mid_at_place"]) / order["mid_at_place"] * 10000 \ + if order["mid_at_place"] > 0 else 0 + if price_moved > 0 and dist_bps > 0: + approach_ratio = price_moved / dist_bps + if approach_ratio < 0.3 and lifetime > 0.5: + is_spoof = True + reasons.append("cancel_before_price_approach") + + # Rule 3: Order size much larger than typical, never fills + if order["sz"] > avg_size * 5 and lifetime < 2.0: + is_spoof = True + reasons.append("oversized_short_lived") + + if is_spoof: + self._spoof_count += 1 + + self._cancel_events.append({ + "spoof": is_spoof, + "time": timestamp, + "lifetime": round(lifetime, 3), + "dist_bps": round(dist_bps, 1), + "reasons": reasons, + }) + + return is_spoof + + def record_fill(self, order_id: str, timestamp: float): + """Record a fill — removes order from tracking, not a spoof.""" + self._orders.pop(order_id, None) + + # ── Metrics ────────────────────────────────────────────── + + def spoof_probability(self) -> float: + """Probability that the current market is being spoofed (0-1). + + Based on recent cancel event ratio and pattern clustering. + """ + recent = [e for e in self._cancel_events + if e["time"] > (self._cancel_events[-1]["time"] if self._cancel_events else 0) - self._window] + if not recent: + return 0.0 + + spoof_recent = sum(1 for e in recent if e["spoof"]) + ratio = spoof_recent / len(recent) + + return min(1.0, ratio * 2.0) # amplify: 50% spoof rate = 100% probability + + def cancel_to_fill_ratio(self) -> float: + """Ratio of cancellations to fills. High ratio = suspicious.""" + total_fills = len(self._fill_events) + if total_fills == 0: + return 1.0 if self._total_cancels > 0 else 0.0 + return self._total_cancels / total_fills + + def spoof_count(self) -> int: + return self._spoof_count + + def summary(self) -> dict: + return { + "spoof_probability": round(self.spoof_probability(), 4), + "spoof_count": self._spoof_count, + "total_orders": self._total_orders, + "total_cancels": self._total_cancels, + "cancel_fill_ratio": round(self.cancel_to_fill_ratio(), 2), + "active_orders": len(self._orders), + } diff --git a/tests/test_advanced_monitors.py b/tests/test_advanced_monitors.py new file mode 100644 index 0000000..039a4b6 --- /dev/null +++ b/tests/test_advanced_monitors.py @@ -0,0 +1,170 @@ +""" +Tests for live/monitors/term_structure.py, liq_waterfall.py, +and microstructure/spoof_detector.py. +""" +from live.monitors.term_structure import TermStructureMonitor +from live.monitors.liq_waterfall import LiquidationWaterfall +from microstructure.spoof_detector import SpoofDetector + + +class TestTermStructure: + def test_initial_no_signal(self): + tsm = TermStructureMonitor() + s = tsm.signal("BTC") + assert s["primary_signal"] == "none" + + def test_basis_computation(self): + tsm = TermStructureMonitor(basis_window=10) + for i in range(10): + tsm.update_perp("BTC", 64500 + i * 10, 0.00001) + tsm.update_quarterly("BTC", 64550 + i * 10) + basis = tsm.quarterly_perp_basis("BTC") + assert basis is not None + assert basis["current_bps"] > 0 # quarterly > perp + assert basis["n_samples"] >= 2 + + def test_zscore_signal(self): + tsm = TermStructureMonitor(basis_window=20) + # Create converging basis (quarterly premium dropping) + for i in range(20): + tsm.update_perp("BTC", 64500, 0.00001) + quarterly_premium = 100 * (1 - i / 20) # declining premium + tsm.update_quarterly("BTC", 64500 + quarterly_premium) + basis = tsm.quarterly_perp_basis("BTC") + assert basis is not None + assert basis["z_score"] < 0 # basis declining below mean + + def test_fair_quarterly_price(self): + tsm = TermStructureMonitor(basis_window=10) + for _ in range(10): + tsm.update_perp("BTC", 64500, 0.00001) + fair = tsm.fair_quarterly_price("BTC") + assert fair is not None + assert fair["perp_price"] == 64500 + assert fair["fair_quarterly"] >= 64500 # positive carry + + def test_signal_on_anomaly(self): + tsm = TermStructureMonitor(basis_window=30) + for _ in range(15): + tsm.update_perp("BTC", 64500, 0.00001) + tsm.update_quarterly("BTC", 64550) + # Spike: quarterly jumps way above fair (>3 sigma) + for _ in range(15): + tsm.update_perp("BTC", 64500, 0.00001) + tsm.update_quarterly("BTC", 65200) # massive premium ~108 bps + s = tsm.signal("BTC") + basis = tsm.quarterly_perp_basis("BTC") + assert basis is not None + assert abs(basis["z_score"]) > 0 # deviation exists + + +class TestLiquidationWaterfall: + def test_initial_no_risk(self): + lw = LiquidationWaterfall() + assert len(lw.at_risk_accounts()) == 0 + + def test_margin_ratio_computation(self): + lw = LiquidationWaterfall() + lw.update_account("0xabc123", [ + {"coin": "BTC", "szi": "1.0", "entryPx": "64000"}, + {"coin": "ETH", "szi": "-10.0", "entryPx": "3100"}, + ], margin_balance=2500) # lower balance → at risk + + risky = lw.at_risk_accounts() + assert len(risky) == 1 + assert risky[0]["level"] in ("danger", "critical") + + def test_safe_account_not_flagged(self): + lw = LiquidationWaterfall() + lw.update_account("0xsafe", [ + {"coin": "BTC", "szi": "0.1", "entryPx": "64000"}, + ], margin_balance=100000) + assert len(lw.at_risk_accounts()) == 0 + + def test_liquidation_order_prediction(self): + lw = LiquidationWaterfall() + lw.update_account("0xwhale", [ + {"coin": "BTC", "szi": "5.0", "entryPx": "64000"}, + {"coin": "ETH", "szi": "-50.0", "entryPx": "3100"}, + {"coin": "SOL", "szi": "1000.0", "entryPx": "140"}, + ], margin_balance=30000) + + # Set book depths: BTC very liquid, SOL very thin + lw.update_book_depth("BTC", 1000000, 1000000) + lw.update_book_depth("ETH", 500000, 500000) + lw.update_book_depth("SOL", 10000, 10000) + + order = lw.predict_liquidation_order("0xwhale") + assert len(order) == 3 + # SOL should be first (thin book, high mm/book ratio) + assert order[0]["predicted_first"] + assert order[0]["coin"] == "SOL" + + def test_signal_when_at_risk(self): + lw = LiquidationWaterfall(danger_margin_ratio=5.0) + lw.update_account("0xrisk", [ + {"coin": "BTC", "szi": "1.0", "entryPx": "64000"}, + ], margin_balance=2000) + lw.update_book_depth("BTC", 10000, 10000) + s = lw.signal("0xrisk") + assert s["action"] == "position" + assert s["liquidation_target"] == "BTC" + + def test_signal_ignores_safe_account(self): + lw = LiquidationWaterfall() + lw.update_account("0xsafe", [ + {"coin": "BTC", "szi": "0.1", "entryPx": "64000"}, + ], margin_balance=100000) + s = lw.signal("0xsafe") + assert s["action"] == "none" + + +class TestSpoofDetector: + def test_initial_probability_zero(self): + sd = SpoofDetector() + assert sd.spoof_probability() == 0.0 + + def test_normal_order_not_spoof(self): + sd = SpoofDetector() + sd.record_place("o1", "bid", 64400, 0.001, 64500, 100.0) + is_spoof = sd.record_cancel("o1", 64500, 105.0) + assert not is_spoof # small order, close to mid, reasonable lifetime + + def test_large_far_quick_cancel_is_spoof(self): + sd = SpoofDetector(size_multiple=2.0) + for _ in range(10): + sd.record_place(f"fill_{_}", "bid", 64400, 0.001, 64500, 0.0) + sd.record_fill(f"fill_{_}", 1.0) + # Large order far from mid, cancelled immediately + sd.record_place("spoof1", "bid", 63000, 10.0, 64500, 200.0) # 1500 bps from mid + is_spoof = sd.record_cancel("spoof1", 64500, 200.1) + assert is_spoof + + def test_oversized_short_lived_is_spoof(self): + sd = SpoofDetector(size_multiple=2.0) + for _ in range(10): + sd.record_place(f"n{_}", "bid", 64400, 0.001, 64500, 0.0) + sd.record_fill(f"n{_}", 1.0) + sd.record_place("big1", "bid", 64400, 50.0, 64500, 300.0) + is_spoof = sd.record_cancel("big1", 64500, 301.5) + assert is_spoof # 50x typical size, <2s lifetime + + def test_spoof_probability_increases(self): + sd = SpoofDetector(window_seconds=5.0) + for _ in range(10): + sd.record_place(f"n{_}", "bid", 64400, 0.001, 64500, 0.0) + sd.record_fill(f"n{_}", 1.0) + # Inject spoofs + for i in range(5): + sd.record_place(f"s{i}", "bid", 63000, 10.0, 64500, 100.0 + i * 0.1) + sd.record_cancel(f"s{i}", 64500, 100.1 + i * 0.1) + assert sd.spoof_probability() > 0 + + def test_summary(self): + sd = SpoofDetector() + sd.record_place("o1", "bid", 64400, 0.001, 64500, 100.0) + sd.record_cancel("o1", 64500, 105.0) + s = sd.summary() + assert "spoof_probability" in s + assert "total_orders" in s + assert s["total_orders"] == 1 diff --git a/tests/test_funding_whipsaw.py b/tests/test_funding_whipsaw.py new file mode 100644 index 0000000..722c77b --- /dev/null +++ b/tests/test_funding_whipsaw.py @@ -0,0 +1,91 @@ +""" +Tests for live/strategies/funding_whipsaw.py — premium index decay trading. +""" +from unittest.mock import patch, MagicMock +from live.strategies.funding_whipsaw import FundingWhipsawTrader + + +class TestFundingWhipsaw: + def test_initial_no_signal_outside_window(self): + """With default 60s entry window and >60s to funding, no signal.""" + trader = FundingWhipsawTrader() + trader.update(mark_px=64500, oracle_px=64480) + with patch.object(trader, 'seconds_to_funding', return_value=300.0): + s = trader.signal() + assert s["action"] == "none" + assert s["reason"] == "outside_entry_window" + + def test_enter_long_on_negative_premium(self): + trader = FundingWhipsawTrader(min_premium_bps=0.5) + trader.update(mark_px=64400, oracle_px=64500) # negative premium: -1.55 bps + with patch.object(trader, 'seconds_to_funding', return_value=30.0): + s = trader.signal() + assert s["action"] == "enter_long" + assert s["side"] == "buy" + assert s["confidence"] > 0 + assert s["size"] > 0 + + def test_enter_short_on_positive_premium(self): + trader = FundingWhipsawTrader(min_premium_bps=0.5) + trader.update(mark_px=64600, oracle_px=64500) # positive premium: +1.55 bps + with patch.object(trader, 'seconds_to_funding', return_value=30.0): + s = trader.signal() + assert s["action"] == "enter_short" + assert s["side"] == "sell" + assert s["confidence"] > 0 + + def test_no_signal_on_small_premium(self): + trader = FundingWhipsawTrader(min_premium_bps=5.0) + trader.update(mark_px=64501, oracle_px=64500) # tiny premium: 0.015 bps + with patch.object(trader, 'seconds_to_funding', return_value=30.0): + s = trader.signal() + assert s["action"] == "none" + assert s["reason"] == "premium_too_small" + + def test_close_after_funding(self): + """After entering, close when funding settles.""" + trader = FundingWhipsawTrader() + trader.update(mark_px=64600, oracle_px=64500) + with patch.object(trader, 'seconds_to_funding', return_value=30.0): + s = trader.signal() + assert s["action"].startswith("enter") + + # Now funding just happened + with patch.object(trader, 'seconds_to_funding', return_value=8*3600 - 2.0): + with patch.object(trader, 'seconds_since_funding', return_value=2.0): + s = trader.signal() + assert s["action"] == "close" + + def test_hold_after_entry(self): + trader = FundingWhipsawTrader() + trader.update(mark_px=64600, oracle_px=64500) + with patch.object(trader, 'seconds_to_funding', return_value=30.0): + trader.signal() # enter + # Next tick, still before funding + with patch.object(trader, 'seconds_to_funding', return_value=25.0): + s = trader.signal() + assert s["action"] == "hold" + + def test_larger_position_with_more_premium(self): + trader = FundingWhipsawTrader(min_premium_bps=0.5, max_size=0.001) + trader.update(mark_px=65100, oracle_px=64500) # large premium + with patch.object(trader, 'seconds_to_funding', return_value=30.0): + s = trader.signal() + assert s["confidence"] > 0.5 # high confidence + assert s["size"] > 0 + + def test_seconds_to_funding_returns_positive(self): + trader = FundingWhipsawTrader() + secs = trader.seconds_to_funding() + assert secs > 0 + assert secs <= 8 * 3600 # Max 8 hours + + def test_signal_includes_premium_info(self): + trader = FundingWhipsawTrader() + trader.update(mark_px=64600, oracle_px=64500) + with patch.object(trader, 'seconds_to_funding', return_value=45.0): + s = trader.signal() + assert "premium_bps" in s + assert "seconds_to_funding" in s + assert "confidence" in s + assert s["premium_bps"] > 0 diff --git a/tests/test_hawkes.py b/tests/test_hawkes.py new file mode 100644 index 0000000..e012b3a --- /dev/null +++ b/tests/test_hawkes.py @@ -0,0 +1,149 @@ +""" +Tests for microstructure/hawkes.py — multivariate Hawkes process calibrator. +""" +import numpy as np +from microstructure.hawkes import ( + HawkesCalibrator, + hawkes_log_likelihood, + hawkes_intensity, + generate_hawkes_events, +) + + +class TestHawkesCalibrator: + + def test_initial_state(self): + cal = HawkesCalibrator(n_dimensions=3) + assert cal.n_dim == 3 + assert cal.mu.shape == (3,) + assert cal.alpha.shape == (3, 3) + assert cal.beta > 0 + + def test_calibrate_on_synthetic_data(self): + """Calibrate on synthetic events from known parameters, check recovery.""" + np.random.seed(42) + # Generate events with known parameters: 3 types + true_mu = np.array([0.5, 0.3, 0.2]) + true_alpha = np.array([ + [0.1, 0.05, 0.02], + [0.03, 0.08, 0.01], + [0.01, 0.02, 0.06], + ]) + true_beta = 2.0 + max_time = 500.0 + + events = generate_hawkes_events(true_mu, true_alpha, true_beta, max_time, seed=42) + assert len(events) >= 3, f"Expected events across 3 types, got {len(events)}" + + cal = HawkesCalibrator(n_dimensions=3) + cal.calibrate(events, max_time) + + # Check that calibrated parameters are within reasonable range + assert np.all(cal.mu > 0), f"mu should be positive, got {cal.mu}" + assert np.all(cal.alpha >= 0), f"alpha should be non-negative, got {cal.alpha}" + assert cal.beta > 0 + + def test_intensity_interpolation(self): + """Intensity should recover to baseline between events and spike after.""" + cal = HawkesCalibrator(n_dimensions=3) + cal.mu = np.array([0.5, 0.3, 0.2]) + cal.alpha = np.array([[0.1, 0, 0], [0, 0, 0], [0, 0, 0]]) + cal.beta = 2.0 + + # Before first event at t=0, intensity = mu (baseline) + i0 = cal.intensity(0, event_history=[], current_time=0.0) + assert abs(i0 - cal.mu[0]) < 0.001 + + # After an event of type 0 at t=0, type 0 intensity should spike + i_after = cal.intensity(0, event_history=[(0, 0.0)], current_time=0.01) + assert i_after > cal.mu[0] + + # After decay, should approach baseline + i_later = cal.intensity(0, event_history=[(0, 0.0)], current_time=5.0) + assert abs(i_later - cal.mu[0]) < 0.05 + + def test_branching_ratio(self): + """Branching ratio should be between 0 and 1.""" + cal = HawkesCalibrator(n_dimensions=3) + cal.mu = np.array([0.5, 0.3, 0.2]) + cal.alpha = np.array([[0.1, 0, 0], [0, 0.1, 0], [0, 0, 0.1]]) + cal.beta = 2.0 + ratio = cal.branching_ratio() + assert 0.0 <= ratio <= 1.0 + + def test_forecast_activity(self): + """Forecast event count in next window.""" + cal = HawkesCalibrator(n_dimensions=3) + cal.mu = np.array([1.0, 0.5, 0.3]) + cal.alpha = np.array([[0.1, 0, 0], [0, 0, 0], [0, 0, 0]]) + cal.beta = 2.0 + + events = [(0, 0.0), (0, 0.5), (0, 1.0), (1, 1.5)] + forecast = cal.forecast_activity(events, current_time=2.0, horizon=5.0) + assert len(forecast) == 3 + assert np.all(forecast >= 0) + + def test_cross_excitation_detected(self): + """Alpha matrix should capture cross-excitation between types.""" + np.random.seed(123) + true_mu = np.array([0.5, 0.3, 0.2]) + true_alpha = np.array([ + [0.2, 0.0, 0.0], + [0.1, 0.1, 0.0], # type 0 excites type 1 + [0.0, 0.0, 0.1], + ]) + true_beta = 3.0 + + events = generate_hawkes_events(true_mu, true_alpha, true_beta, 500.0, seed=123) + cal = HawkesCalibrator(n_dimensions=3) + cal.calibrate(events, 500.0) + + # Type 0 should excite type 1: alpha[1,0] > 0 + assert cal.alpha[1, 0] >= 0, f"Expected cross-excitation, got alpha[1,0]={cal.alpha[1,0]}" + + +class TestHawkesFunctions: + + def test_log_likelihood_improves_with_fit(self): + np.random.seed(99) + true_mu = np.array([0.5, 0.3, 0.2]) + true_alpha = np.array([[0.15, 0.0, 0.0], [0.0, 0.1, 0.0], [0.0, 0.0, 0.05]]) + true_beta = 2.5 + + events = generate_hawkes_events(true_mu, true_alpha, true_beta, 200.0, seed=99) + + # Bad parameters + bad_mu = np.array([0.1, 0.1, 0.1]) + bad_alpha = np.zeros((3, 3)) + ll_bad = hawkes_log_likelihood(events, bad_mu, bad_alpha, true_beta, 200.0) + + # Good parameters (close to truth) + ll_good = hawkes_log_likelihood(events, true_mu, true_alpha, true_beta, 200.0) + + assert ll_good > ll_bad, f"Good params should give higher likelihood: {ll_good} vs {ll_bad}" + + def test_intensity_function(self): + mu = np.array([0.5, 0.3]) + alpha = np.array([[0.1, 0.0], [0.0, 0.05]]) + beta = 2.0 + events = [(0, 0.0), (0, 0.5), (1, 1.0)] + + intensity = hawkes_intensity(mu, alpha, beta, events, current_time=0.6, dim=0) + assert intensity > mu[0], f"After events, intensity should exceed baseline" + + def test_generate_events_produces_timestamps(self): + np.random.seed(42) + events = generate_hawkes_events( + np.array([0.5, 0.3]), + np.array([[0.1, 0.0], [0.0, 0.05]]), + beta=2.0, + max_time=100.0, + seed=42, + ) + assert len(events) > 0 + # Each event should be (type, time) + for ev in events: + assert isinstance(ev, tuple) + assert len(ev) == 2 + assert ev[0] in (0, 1) + assert 0 <= ev[1] <= 100.0 diff --git a/tests/test_hlp_vault.py b/tests/test_hlp_vault.py new file mode 100644 index 0000000..690c84f --- /dev/null +++ b/tests/test_hlp_vault.py @@ -0,0 +1,153 @@ +""" +Tests for live/monitors/hlp_vault.py — HLP protocol-level market maker tracking. +""" +from unittest.mock import patch +from live.monitors.hlp_vault import HlpVaultMonitor + +HLP_ADDRESS = "0xfefefefefefefefefefefefefefefefefefefefe" + +def _mock_meta(): return {"universe": [ + {"name": "BTC", "szDecimals": 5}, + {"name": "ETH", "szDecimals": 6}, + {"name": "SOL", "szDecimals": 7}, +]} + +def _mock_asset_ctxs(): return [ + {"funding": "0.00001", "markPx": "64500", "oraclePx": "64480", "openInterest": "50000000"}, + {"funding": "0.000005", "markPx": "3200", "oraclePx": "3195", "openInterest": "30000000"}, + {"funding": "0.00002", "markPx": "140", "oraclePx": "139.5", "openInterest": "10000000"}, +] + +def _mock_clearinghouse(positions=None): + aps = [] + if positions: + for coin, (side, szi, entry_px, upnl) in positions.items(): + aps.append({"type": "oneWay", "position": { + "coin": coin, "side": side, "szi": str(szi), + "entryPx": str(entry_px), "unrealizedPnl": str(upnl), + }}) + return {"assetPositions": aps, "withdrawable": "1000000"} + + +class TestHlpVaultMonitor: + + def test_initial_state_empty(self): + monitor = HlpVaultMonitor(testnet=True) + s = monitor.summary() + assert s["assets_tracked"] == 0 + assert s["total_delta_usd"] == 0.0 + + def test_update_populates_positions(self): + monitor = HlpVaultMonitor(testnet=True) + with patch.object(monitor, '_api_post') as m: + m.side_effect = [ + [_mock_meta(), _mock_asset_ctxs()], # metaAndAssetCtxs → list + _mock_clearinghouse({"BTC": ("A", 10.5, 64000, 5250), + "ETH": ("B", 50.0, 3100, -2500)}), # clearinghouseState → dict + ] + monitor.update() + assert monitor.position("BTC") < 0 # side=A = short + assert monitor.position("ETH") > 0 # side=B = long + assert monitor.summary()["assets_tracked"] >= 2 + + def test_delta_exposure_usd(self): + monitor = HlpVaultMonitor(testnet=True) + with patch.object(monitor, '_api_post') as m: + m.side_effect = [ + [_mock_meta(), _mock_asset_ctxs()], + _mock_clearinghouse({"BTC": ("A", 10.0, 64000, 5000)}), + ] + monitor.update() + delta = monitor.delta_exposure() + assert "BTC" in delta + assert abs(delta["BTC"]["notional_usd"]) > 600000 + + def test_is_overextended(self): + monitor = HlpVaultMonitor(testnet=True, overextended_threshold=5.0) + with patch.object(monitor, '_api_post') as m: + m.side_effect = [ + [_mock_meta(), _mock_asset_ctxs()], + _mock_clearinghouse({"BTC": ("A", 100.0, 64000, 50000)}), + ] + monitor.update() + assert monitor.is_overextended("BTC") + + def test_not_overextended_with_small_position(self): + monitor = HlpVaultMonitor(testnet=True, overextended_threshold=5.0) + with patch.object(monitor, '_api_post') as m: + m.side_effect = [ + [_mock_meta(), _mock_asset_ctxs()], + _mock_clearinghouse({"BTC": ("A", 1.0, 64000, 500)}), + ] + monitor.update() + assert not monitor.is_overextended("BTC") + + def test_rebalancing_signal_long(self): + monitor = HlpVaultMonitor(testnet=True, overextended_threshold=5.0) + with patch.object(monitor, '_api_post') as m: + m.side_effect = [ + [_mock_meta(), _mock_asset_ctxs()], + _mock_clearinghouse({"BTC": ("A", 100.0, 64000, 50000)}), + ] + monitor.update() + signal = monitor.rebalancing_signal("BTC") + assert signal["overextended"] + assert signal["signal"] in ("fade_short", "fade_long", "neutral") + + def test_rebalancing_signal_short(self): + monitor = HlpVaultMonitor(testnet=True, overextended_threshold=5.0) + with patch.object(monitor, '_api_post') as m: + m.side_effect = [ + [_mock_meta(), _mock_asset_ctxs()], + _mock_clearinghouse({"ETH": ("B", 2000.0, 3100, 50000)}), + ] + monitor.update() + signal = monitor.rebalancing_signal("ETH") + assert signal["overextended"] + + def test_toxicity_score_zero_with_no_data(self): + monitor = HlpVaultMonitor(testnet=True) + assert monitor.toxicity_score() == 0.0 + + def test_toxicity_score_detects_losing_flow(self): + monitor = HlpVaultMonitor(testnet=True) + with patch.object(monitor, '_api_post') as m: + m.side_effect = [ + [_mock_meta(), _mock_asset_ctxs()], + _mock_clearinghouse({ + "BTC": ("A", 10.0, 64000, -50000), + "ETH": ("A", 50.0, 3100, -25000), + }), + ] + monitor.update() + score = monitor.toxicity_score() + assert score > 0 + + def test_historical_tracking(self): + monitor = HlpVaultMonitor(testnet=True) + with patch.object(monitor, '_api_post') as m: + m.side_effect = [ + [_mock_meta(), _mock_asset_ctxs()], + _mock_clearinghouse({"BTC": ("A", 10.0, 64000, 0)}), + [_mock_meta(), _mock_asset_ctxs()], + _mock_clearinghouse({"BTC": ("A", 12.0, 64100, 1000)}), + ] + monitor.update() + monitor.update() + history = monitor.delta_history("BTC") + assert len(history) == 2 + + def test_summary_includes_all_fields(self): + monitor = HlpVaultMonitor(testnet=True) + with patch.object(monitor, '_api_post') as m: + m.side_effect = [ + [_mock_meta(), _mock_asset_ctxs()], + _mock_clearinghouse({"BTC": ("A", 10.0, 64000, 5000)}), + ] + monitor.update() + s = monitor.summary() + assert "assets_tracked" in s + assert "total_delta_usd" in s + assert "toxicity_score" in s + assert "overextended_assets" in s + assert "signals" in s