diff --git a/live/monitors/sequencer_latency.py b/live/monitors/sequencer_latency.py new file mode 100644 index 0000000..d556350 --- /dev/null +++ b/live/monitors/sequencer_latency.py @@ -0,0 +1,181 @@ +""" +Sequencer latency & stale-state arbitrage detector. + +Hyperliquid uses a centralized sequencer for off-chain order matching. +There is a microsecond-to-millisecond delta between: + 1. WebSocket trade print (fastest) + 2. REST API state (slower) + 3. Cross-margin collateral recalculation (slowest) + +When a liquidation hits one asset, cross-margin collateral drops for ALL +assets in that portfolio — but the sequencer may not have updated margin +limits on OTHER order books yet. This creates a stale-state window. + +Strategy: detect liquidations on BTC, then hit ETH book before margins update. +""" + +from __future__ import annotations + +import time +from collections import deque +from typing import Optional + + +class SequencerLatencyDetector: + """Detect and exploit sequencer state-update latency. + + Tracks WebSocket vs REST timestamps to measure the gap, + and identifies stale-state windows after liquidation events. + """ + + def __init__( + self, + window_seconds: float = 60.0, + stale_threshold_ms: float = 50.0, # >50ms REST lag = stale + max_history: int = 1000, + ): + self._window = window_seconds + self._stale_threshold = stale_threshold_ms + self._max_history = max_history + + # Latency tracking per channel per coin + self._ws_timestamps: dict[str, dict[str, deque[float]]] = {} + self._rest_timestamps: dict[str, dict[str, deque[float]]] = {} + self._latency_measurements: deque = deque(maxlen=max_history) + + # Liquidation event log + self._liquidations: deque = deque(maxlen=500) + + # Sequencer state + self._stale_window_active: bool = False + self._stale_assets: list[str] = [] + self._stale_start: float = 0.0 + + # ── Data feed ──────────────────────────────────────────── + + def record_ws_event(self, coin: str, channel: str, exchange_ts_ms: int): + """Record WebSocket event timestamp (fastest).""" + local_ts = time.time() * 1000 + c = coin.upper() + ch = channel + + if c not in self._ws_timestamps: + self._ws_timestamps[c] = {} + self._ws_timestamps[c].setdefault(ch, deque(maxlen=self._max_history)).append(exchange_ts_ms) + + # Measure latency: how much time between exchange and local receipt + lat = local_ts - exchange_ts_ms + self._latency_measurements.append({ + "time": time.time(), + "coin": c, + "channel": ch, + "latency_ms": round(lat, 2), + "direction": "ws", + }) + + def record_rest_event(self, coin: str, endpoint: str, exchange_ts_ms: int): + """Record REST API timestamp (slower).""" + c = coin.upper() + if c not in self._rest_timestamps: + self._rest_timestamps[c] = {} + self._rest_timestamps[c].setdefault(endpoint, deque(maxlen=self._max_history)).append(exchange_ts_ms) + + def record_liquidation(self, coin: str, size: float, price: float, timestamp_ms: int): + """Record a liquidation event — triggers stale-state analysis.""" + self._liquidations.append({ + "time": time.time(), + "coin": coin.upper(), + "size": size, + "price": price, + "exchange_ts": timestamp_ms, + }) + + # ── Analysis ───────────────────────────────────────────── + + def ws_rest_latency(self, coin: str, channel: str = "l2book") -> dict | None: + """Measure latency between WebSocket and REST for a coin/channel. + + Returns median latency (ms), count of measurements, and stale flag. + """ + ws = list(self._ws_timestamps.get(coin.upper(), {}).get(channel, [])) + rest = list(self._rest_timestamps.get(coin.upper(), {}).get(channel, [])) + if not ws or not rest: + return None + + # Compare latest timestamps — if WS is newer than REST by > threshold + ws_latest = ws[-1] + rest_latest = rest[-1] + delta = rest_latest - ws_latest + + return { + "coin": coin.upper(), + "channel": channel, + "ws_latest_ms": ws_latest, + "rest_latest_ms": rest_latest, + "delta_ms": round(delta, 2), + "is_stale": delta > self._stale_threshold, + "measurements": min(len(ws), len(rest)), + } + + def avg_transport_latency(self) -> dict: + """Average WebSocket transport latency (exchange → local). + + Returns p50, p99, max, and sample count. + """ + lats = [m["latency_ms"] for m in self._latency_measurements] + if not lats: + return {"p50_ms": 0, "p99_ms": 0, "max_ms": 0, "count": 0} + + sorted_lats = sorted(lats) + n = len(sorted_lats) + return { + "p50_ms": round(sorted_lats[int(n * 0.50)], 2), + "p99_ms": round(sorted_lats[int(n * 0.99)], 2), + "max_ms": round(max(lats), 2), + "count": n, + } + + def recent_liquidations(self) -> list[dict]: + """Liquidation events in the last window_seconds.""" + cutoff = time.time() - self._window + return [l for l in self._liquidations if l["time"] >= cutoff] + + def stale_state_signal(self) -> dict: + """Detect if a stale-state window is active after liquidation. + + If a liquidation just occurred and REST lag is > threshold: + - Mark affected assets as stale + - Signal which other assets may have outdated margin limits + """ + recent_liqs = self.recent_liquidations() + if not recent_liqs: + return {"stale_window_active": False, "signal": "none"} + + latest_liq = recent_liqs[-1] + age_ms = (time.time() - latest_liq["time"]) * 1000 + + # Check if any asset has REST lag exceeding threshold + stale_assets = [] + for coin, channels in self._ws_timestamps.items(): + for channel in channels: + lat = self.ws_rest_latency(coin, channel) + if lat and lat["is_stale"]: + stale_assets.append(lat) + + is_stale = len(stale_assets) > 0 and age_ms < 1000 # within 1s of liquidation + + return { + "stale_window_active": is_stale, + "signal": "stale_state_detected" if is_stale else "none", + "liquidation_coin": latest_liq["coin"], + "liquidation_age_ms": round(age_ms, 2), + "stale_assets": [s["coin"] for s in stale_assets], + "recommended_action": "check_margin_limits_on_STALE_ASSETS" if is_stale else "none", + } + + def summary(self) -> dict: + return { + "avg_latency": self.avg_transport_latency(), + "recent_liquidations": len(self.recent_liquidations()), + "stale_state": self.stale_state_signal(), + } diff --git a/live/monitors/triangular_arb.py b/live/monitors/triangular_arb.py new file mode 100644 index 0000000..af6f7f6 --- /dev/null +++ b/live/monitors/triangular_arb.py @@ -0,0 +1,213 @@ +""" +Cross-exchange triangular latency arbitrage detector. + +Latency arb isn't just A vs B. It's A → B → C across three venues. +If BTC/USDT is slow to update on Venue 1, but Venue 2 is fast on +ETH/BTC and Venue 3 is fast on ETH/USDT, the slow leg creates a +triangular arbitrage opportunity. + +Path: Buy BTC/USDT on slow venue → Sell ETH/BTC on fast venue → Sell ETH/USDT on fast venue + +Strategy: monitor 3-venue latency simultaneously, detect when one leg +lags, compute implied arbitrage spread, and signal execution-ready +opportunities. +""" + +from __future__ import annotations + +from collections import deque +from typing import Optional + + +class TriangularLatencyArb: + """Detect triangular arbitrage from cross-venue latency discrepancies. + + Monitors up to 3 venues with configurable latency estimates. + When one venue lags on a specific pair, the triangle becomes profitable. + + Usage: + arb = TriangularLatencyArb() + arb.update_price("hl", "BTC-USDT", 64500, latency_ms=5) + arb.update_price("hl", "ETH-BTC", 0.049, latency_ms=8) + arb.update_price("hl", "ETH-USDT", 3160, latency_ms=6) + opportunities = arb.detect() + """ + + def __init__( + self, + min_spread_bps: float = 0.5, # minimum bps profit to signal + max_latency_diff_ms: float = 200.0, # max venue latency gap to consider + window: int = 50, + ): + self._min_spread = min_spread_bps + self._max_latency = max_latency_diff_ms + self._window = window + + # venue → pair → (price, latency_ms, timestamp) + self._prices: dict[str, dict[str, deque[tuple[float, float, float]]]] = {} + self._latencies: dict[str, float] = {} # venue → avg latency + + # ── Data feed ──────────────────────────────────────────── + + def update_price(self, venue: str, pair: str, price: float, latency_ms: float = 0): + """Record a price snapshot from a venue with its latency.""" + v = venue.lower() + p = pair.upper().replace("/", "-") + + self._prices.setdefault(v, {}).setdefault( + p, deque(maxlen=self._window) + ).append((price, latency_ms, latency_ms)) # (price, latency, timestamp_simple) + + # Update venue latency estimate (EMA) + old_lat = self._latencies.get(v, latency_ms) + self._latencies[v] = old_lat * 0.9 + latency_ms * 0.1 + + def get_price(self, venue: str, pair: str) -> float | None: + """Get latest price from a venue.""" + dq = self._prices.get(venue.lower(), {}).get(pair.upper().replace("/", "-")) + return dq[-1][0] if dq else None + + def get_latency(self, venue: str) -> float: + return self._latencies.get(venue.lower(), 0) + + # ── Arbitrage detection ────────────────────────────────── + + def detect(self) -> dict: + """Detect triangular arbitrage opportunities across venues. + + Triangle: BTC-USDT → ETH-BTC → ETH-USDT + If any leg is on a slower venue, implied profit exists. + + Returns list of opportunities with profit, confidence, and execution plan. + """ + opportunities = [] + pairs = ["BTC-USDT", "ETH-BTC", "ETH-USDT"] + venues = list(self._prices.keys()) + + if len(venues) < 1: + return {"opportunities": [], "venue_count": len(venues)} + + # Find latency gaps between venues + if len(venues) >= 2: + lat_gaps = self._latency_gaps() + else: + lat_gaps = {} + + # For each venue, compute implied cross-rate vs direct + for v in venues: + btc_usdt = self.get_price(v, "BTC-USDT") + eth_btc = self.get_price(v, "ETH-BTC") + eth_usdt = self.get_price(v, "ETH-USDT") + + if not (btc_usdt and eth_btc and eth_usdt): + continue + + # Implied ETH-USDT from triangle: BTC-USDT × ETH-BTC + implied_eth = btc_usdt * eth_btc + implied_bps = (eth_usdt - implied_eth) / implied_eth * 10000 + + if abs(implied_bps) > self._min_spread: + opportunities.append({ + "venue": v, + "type": "internal_triangular", + "btc_usdt": btc_usdt, + "eth_btc": eth_btc, + "eth_usdt": eth_usdt, + "implied_eth_usdt": round(implied_eth, 2), + "spread_bps": round(implied_bps, 2), + "direction": "sell_eth" if implied_bps > 0 else "buy_eth", + "latency_ms": round(self._latencies.get(v, 0), 2), + }) + + # Cross-venue: check if one venue's slow leg creates arb + if len(venues) >= 2: + for i, v1 in enumerate(venues): + for v2 in venues[i + 1:]: + opp = self._cross_venue_arb(v1, v2, pairs) + if opp: + opportunities.append(opp) + + return { + "opportunities": opportunities[:10], + "venue_count": len(venues), + "latency_gaps": lat_gaps, + "best_opportunity": max(opportunities, key=lambda o: abs(o["spread_bps"])) if opportunities else None, + } + + def _cross_venue_arb(self, v1: str, v2: str, pairs: list[str]) -> dict | None: + """Check if buying on slow venue, selling on fast venue is profitable.""" + best_opp = None + best_bps = 0 + + for pair in pairs: + p1 = self.get_price(v1, pair) + p2 = self.get_price(v2, pair) + if not p1 or not p2: + continue + + spread_bps = abs(p2 - p1) / p1 * 10000 + lat_diff = abs(self.get_latency(v1) - self.get_latency(v2)) + + if spread_bps > self._min_spread and lat_diff > 10: + opp = { + "type": "cross_venue", + "pair": pair, + "slow_venue": v1 if self.get_latency(v1) > self.get_latency(v2) else v2, + "fast_venue": v1 if self.get_latency(v1) < self.get_latency(v2) else v2, + "buy_at": round(min(p1, p2), 2), + "sell_at": round(max(p1, p2), 2), + "spread_bps": round(spread_bps, 2), + "latency_diff_ms": round(lat_diff, 2), + } + if spread_bps > best_bps: + best_bps = spread_bps + best_opp = opp + + return best_opp + + def _latency_gaps(self) -> dict: + """Compute latency gaps between all venue pairs.""" + venues = list(self._latencies.keys()) + gaps = {} + for i, v1 in enumerate(venues): + for v2 in venues[i + 1:]: + diff = abs(self._latencies.get(v1, 0) - self._latencies.get(v2, 0)) + key = f"{v1}_{v2}" + gaps[key] = round(diff, 2) + return gaps + + def signal(self) -> dict: + """Generate trading signal for cross-venue triangular arb.""" + result = self.detect() + opps = result.get("opportunities", []) + + if not opps: + return {"action": "none", "reason": "no_opportunity"} + + best = opps[0] + if best["type"] == "cross_venue" and abs(best["spread_bps"]) > self._min_spread * 2: + return { + "action": "arbitrage", + "type": "cross_venue", + "details": best, + "confidence": min(1.0, abs(best["spread_bps"]) / (self._min_spread * 5)), + } + elif best["type"] == "internal_triangular" and abs(best["spread_bps"]) > self._min_spread * 3: + return { + "action": "arbitrage", + "type": "triangular", + "details": best, + "confidence": min(1.0, abs(best["spread_bps"]) / (self._min_spread * 5)), + } + + return {"action": "monitor", "reason": "spread_too_small", "best_bps": round(best["spread_bps"], 2)} + + def summary(self) -> dict: + opps = self.detect() + return { + "venues": list(self._prices.keys()), + "latencies": self._latencies, + "opportunities_count": len(opps.get("opportunities", [])), + "best_opportunity": opps.get("best_opportunity"), + "signal": self.signal(), + } diff --git a/microstructure/dealer_gex.py b/microstructure/dealer_gex.py new file mode 100644 index 0000000..abb7a82 --- /dev/null +++ b/microstructure/dealer_gex.py @@ -0,0 +1,179 @@ +""" +Dealer Gamma Exposure (GEX) estimation and pinning/fading signals. + +In options markets, dealers delta-hedge their portfolios. Their gamma +exposure determines whether they amplify or suppress price movements: + - Long gamma → dealers buy low, sell high → suppress vol, "pin" price + - Short gamma → dealers buy high, sell low → amplify vol, create gamma squeeze + +GEX = Σ (gamma_per_contract × open_interest × spot_price) + +For crypto (no options on Hyperliquid yet), we approximate GEX using: + - BTC/ETH options open interest from Deribit + - Estimated dealer gamma via Black-Scholes + - Pin levels at strikes with max GEX + +Strategy: when dealer GEX is heavily positive at a strike, fade breakouts. +When GEX is negative, ride the momentum with the dealers. +""" + +from __future__ import annotations + +import math +from collections import defaultdict +from typing import Optional + + +class DealerGEX: + """Estimate dealer gamma exposure and generate pinning/fading signals. + + Uses option open interest data and Black-Scholes gamma to compute + per-strike GEX. Aggregates to net GEX and identifies pin levels. + + Usage: + gex = DealerGEX(spot=64500, risk_free=0.05) + gex.add_strike(strike=65000, call_oi=500, put_oi=300, expiry_days=7, iv=0.60) + print(gex.pin_levels()) + print(gex.signal()) + """ + + def __init__( + self, + spot: float = 0.0, + risk_free: float = 0.05, + pin_threshold: float = 0.01, # 1% of spot to consider "at pin level" + ): + self._spot = spot + self._rf = risk_free + self._pin_threshold = pin_threshold + + self._strikes: list[dict] = [] # [{strike, call_oi, put_oi, gamma, gex_usd, ...}] + self._net_gex: float = 0.0 + self._pin_levels: list[dict] = [] + + # ── Data feed ──────────────────────────────────────────── + + def set_spot(self, spot: float): + self._spot = spot + + def add_strike( + self, + strike: float, + call_oi: float = 0, + put_oi: float = 0, + expiry_days: float = 30, + iv: float = 0.50, + ): + """Add OI data at a specific strike.""" + gamma_call = self._bs_gamma(strike, "call", expiry_days, iv) + gamma_put = self._bs_gamma(strike, "put", expiry_days, iv) + + # GEX per strike = gamma × OI × spot² × 0.01 (scaling) + gex_call = gamma_call * call_oi * self._spot * self._spot * 0.01 + gex_put = gamma_put * put_oi * self._spot * self._spot * 0.01 + + self._strikes.append({ + "strike": strike, + "call_oi": call_oi, + "put_oi": put_oi, + "gamma_call": round(gamma_call, 8), + "gamma_put": round(gamma_put, 8), + "gex_call_usd": round(gex_call, 2), + "gex_put_usd": round(gex_put, 2), + "gex_total_usd": round(gex_call + gex_put, 2), + }) + + def compute(self): + """Aggregate GEX and find pin levels.""" + self._net_gex = sum(s["gex_total_usd"] for s in self._strikes) + + # Find pin levels: strikes where GEX is highest (magnets) + # and is positive (dealers long gamma = pinning) + if self._strikes: + max_gex = max(abs(s["gex_total_usd"]) for s in self._strikes) + self._pin_levels = [ + {"strike": s["strike"], "gex_usd": s["gex_total_usd"], "is_pin": s["gex_total_usd"] > 0} + for s in self._strikes if abs(s["gex_total_usd"]) > max_gex * 0.3 + ] + self._pin_levels.sort(key=lambda l: abs(l["gex_usd"]), reverse=True) + + # ── Signals ───────────────────────────────────────────── + + def nearest_pin(self) -> dict | None: + """Find the nearest pin level to current spot.""" + if not self._spot or not self._strikes: + return None + + distances = [] + for s in self._strikes: + dist_pct = abs(s["strike"] - self._spot) / self._spot + if s["gex_total_usd"] > 0: # only positive GEX acts as pin + distances.append({ + "strike": s["strike"], + "distance_pct": round(dist_pct * 100, 2), + "gex_usd": s["gex_total_usd"], + "strength": "strong" if s["gex_total_usd"] > abs(self._net_gex) * 0.1 else "weak", + }) + + if not distances: + return None + return min(distances, key=lambda d: d["distance_pct"]) + + def signal(self) -> dict: + """Generate trading signal based on dealer GEX. + + Returns: + signal: "fade_breakout", "ride_momentum", "neutral" + reason: explanation + """ + self.compute() + if not self._spot: + return {"signal": "neutral", "reason": "no_spot"} + + pin = self.nearest_pin() + dist_pct = pin["distance_pct"] if pin else 100.0 + + if self._net_gex > 1e6 and pin and dist_pct < self._pin_threshold * 100: + return { + "signal": "fade_breakout", + "reason": f"Dealers long gamma ${self._net_gex:,.0f} — pinning at ${pin['strike']:,.0f} ({dist_pct:.1f}%)", + "net_gex_usd": round(self._net_gex, 2), + "pin_strike": pin["strike"], + "direction": "mean_reversion", + } + elif self._net_gex < -1e6: + return { + "signal": "ride_momentum", + "reason": f"Dealers short gamma ${self._net_gex:,.0f} — amplifying moves", + "net_gex_usd": round(self._net_gex, 2), + "direction": "trend_following", + } + else: + return {"signal": "neutral", "reason": "balanced_gex", "net_gex_usd": round(self._net_gex, 2)} + + # ── Black-Scholes gamma ───────────────────────────────── + + def _bs_gamma(self, strike: float, opt_type: str, T_days: float, sigma: float) -> float: + """Black-Scholes gamma for a European option.""" + if self._spot <= 0 or strike <= 0 or T_days <= 0 or sigma <= 0: + return 0.0 + + T = T_days / 365.0 + d1 = (math.log(self._spot / strike) + (self._rf + 0.5 * sigma * sigma) * T) / (sigma * math.sqrt(T)) + phi = math.exp(-d1 * d1 / 2.0) / math.sqrt(2.0 * math.pi) + + return phi / (self._spot * sigma * math.sqrt(T)) + + # ── Summary ───────────────────────────────────────────── + + def summary(self) -> dict: + pin = self.nearest_pin() + return { + "spot": self._spot, + "net_gex_usd": round(self._net_gex, 2), + "net_gex_m": round(self._net_gex / 1e6, 2), + "strikes_tracked": len(self._strikes), + "pin_levels": self._pin_levels[:5], + "nearest_pin": pin, + "signal": self.signal(), + } diff --git a/microstructure/tick_regime.py b/microstructure/tick_regime.py new file mode 100644 index 0000000..bee28a7 --- /dev/null +++ b/microstructure/tick_regime.py @@ -0,0 +1,191 @@ +""" +Tick-size and lot-size regime exploitation. + +Exchanges dynamically adjust minimum price increments (tick size) +and minimum order sizes (lot size) based on asset price. When an +asset crosses these thresholds, spreads widen or narrow instantly. + +Hyperliquid tick sizes (example — actual values from exchange): + BTC: 0.1 USD tick (<$100K), 0.5 USD tick ($100K-$1M) — approximate + ETH: 0.01 USD tick + SOL: 0.001 USD tick + +Strategy: when price crosses a tick-size boundary, be the first to adjust +your quoting engine. Competitors quoting at old tick sizes will be either +non-competitive (too wide) or illegal (too tight). Capture the spread gap. +""" + +from __future__ import annotations + +from collections import deque +from typing import Optional + + +# Hyperliquid tick size schedule (example values — confirm from exchange docs) +TICK_SCHEDULE = { + "BTC": [ + (float("-inf"), 100000, 0.1), + (100000, float("inf"), 0.5), + ], + "ETH": [ + (float("-inf"), 10000, 0.01), + (10000, float("inf"), 0.05), + ], + "SOL": [ + (float("-inf"), 1000, 0.001), + (1000, float("inf"), 0.005), + ], + # Default for unknown coins + "DEFAULT": [ + (float("-inf"), float("inf"), 0.01), + ], +} + + +class TickRegimeMonitor: + """Monitor and exploit tick-size regime changes. + + Tracks when an asset price approaches a tick-size boundary + and generates signals to adjust quoting before competitors. + """ + + def __init__( + self, + approach_threshold_pct: float = 1.0, # 1% from boundary + history_window: int = 100, + ): + self._approach_threshold = approach_threshold_pct / 100 + self._prices: dict[str, deque[float]] = {} + self._current_tick: dict[str, float] = {} + self._approaching_boundary: dict[str, dict] = {} + self._window = history_window + + # ── Data feed ──────────────────────────────────────────── + + def get_tick_size(self, coin: str, price: float) -> float: + """Get current tick size for a coin at given price.""" + schedule = TICK_SCHEDULE.get(coin.upper(), TICK_SCHEDULE["DEFAULT"]) + for lo, hi, tick in schedule: + if lo <= price < hi: + return tick + return schedule[-1][2] # fallback + + def update_price(self, coin: str, price: float): + """Feed a new price observation.""" + c = coin.upper() + self._prices.setdefault(c, deque(maxlen=self._window)).append(price) + self._current_tick[c] = self.get_tick_size(c, price) + + # Check if approaching boundary + schedule = TICK_SCHEDULE.get(c, TICK_SCHEDULE["DEFAULT"]) + for lo, hi, tick in schedule: + if lo <= price < hi: + # Check approach to upper boundary + if hi != float("inf"): + dist_up = (hi - price) / price + if 0 < dist_up < self._approach_threshold: + next_tick = self.get_tick_size(c, hi) + if next_tick != tick: + self._approaching_boundary[c] = { + "direction": "up", + "current_tick": tick, + "next_tick": next_tick, + "boundary_price": hi, + "distance_pct": round(dist_up * 100, 2), + "crosses_in_bars": self._estimate_cross_bars(c, hi), + } + return + # Check approach to lower boundary + if lo != float("-inf"): + dist_down = (price - lo) / price + if 0 < dist_down < self._approach_threshold: + prev_tick = self.get_tick_size(c, lo - 1) + if prev_tick != tick: + self._approaching_boundary[c] = { + "direction": "down", + "current_tick": tick, + "next_tick": prev_tick, + "boundary_price": lo, + "distance_pct": round(dist_down * 100, 2), + "crosses_in_bars": self._estimate_cross_bars(c, lo), + } + return + + self._approaching_boundary.pop(c, None) + + def _estimate_cross_bars(self, coin: str, boundary_price: float) -> int: + """Estimate how many bars until price crosses the boundary.""" + prices = list(self._prices.get(coin, [])) + if len(prices) < 5: + return -1 + + recent = prices[-min(20, len(prices)):] + if len(recent) < 5: + return -1 + + # Simple linear trend estimate + xs = list(range(len(recent))) + n = len(xs) + sum_x = sum(xs) + sum_y = sum(recent) + sum_xy = sum(x * y for x, y in zip(xs, recent)) + sum_x2 = sum(x * x for x in xs) + + slope = (n * sum_xy - sum_x * sum_y) / (n * sum_x2 - sum_x * sum_x) if (n * sum_x2 - sum_x * sum_x) != 0 else 0 + last_price = recent[-1] + + if abs(slope) < 1e-10: + return -1 + + bars_to_cross = int((boundary_price - last_price) / slope) + return max(0, bars_to_cross) if bars_to_cross > 0 else -1 + + # ── Signal ────────────────────────────────────────────── + + def signal(self, coin: str) -> dict: + """Generate tick-size regime signal. + + Returns action to take before competitors adjust. + """ + c = coin.upper() + boundary = self._approaching_boundary.get(c) + + current_tick = self._current_tick.get(c, 0.01) + + if not boundary: + return { + "action": "none", + "current_tick": current_tick, + "reason": "no_regime_change", + } + + tick_delta = boundary["next_tick"] - boundary["current_tick"] + bars = boundary["crosses_in_bars"] + + if tick_delta > 0: + # Tick size INCREASING → spread will widen → HFTs will post wider quotes + # Be the first to widen to the new tick + action = "widen_quotes" + reason = f"Tick increasing {boundary['current_tick']}→{boundary['next_tick']} " + reason += f"({boundary['distance_pct']}% away" + (f", ~{bars} bars" if bars > 0 else "") + ")" + else: + # Tick size DECREASING → spread will narrow → HFTs will tighten + # Be the first to tighten to the new tick + action = "tighten_quotes" + reason = f"Tick decreasing {boundary['current_tick']}→{boundary['next_tick']} " + reason += f"({boundary['distance_pct']}% away" + (f", ~{bars} bars" if bars > 0 else "") + ")" + + return { + "action": action, + "current_tick": boundary["current_tick"], + "next_tick": boundary["next_tick"], + "tick_delta": round(tick_delta, 5), + "boundary_price": boundary["boundary_price"], + "distance_pct": boundary["distance_pct"], + "estimated_bars_to_cross": bars, + "reason": reason, + "urgent": boundary["distance_pct"] < 0.5, # within 0.5% = urgent + } + + def summary(self) -> dict: + return {coin: self.signal(coin) for coin in self._prices} diff --git a/tests/test_remaining_strategies.py b/tests/test_remaining_strategies.py new file mode 100644 index 0000000..ebdfc3d --- /dev/null +++ b/tests/test_remaining_strategies.py @@ -0,0 +1,173 @@ +""" +Tests for sequencer latency, dealer GEX, tick regime, triangular arb. +""" +from live.monitors.sequencer_latency import SequencerLatencyDetector +from microstructure.dealer_gex import DealerGEX +from microstructure.tick_regime import TickRegimeMonitor +from live.monitors.triangular_arb import TriangularLatencyArb + + +class TestSequencerLatency: + def test_initial_state(self): + sld = SequencerLatencyDetector() + lat = sld.avg_transport_latency() + assert lat["count"] == 0 + + def test_ws_event_recording(self): + sld = SequencerLatencyDetector() + now = 1705312800000 + sld.record_ws_event("BTC", "l2book", now) + lat = sld.avg_transport_latency() + assert lat["count"] == 1 + + def test_stale_detection(self): + sld = SequencerLatencyDetector(stale_threshold_ms=100) + sld.record_ws_event("BTC", "l2book", 1705312800000) + sld.record_rest_event("BTC", "l2book", 1705312800200) # 200ms later + lat = sld.ws_rest_latency("BTC", "l2book") + assert lat is not None + assert lat["is_stale"] + + def test_not_stale_when_close(self): + sld = SequencerLatencyDetector(stale_threshold_ms=100) + sld.record_ws_event("BTC", "l2book", 1705312800000) + sld.record_rest_event("BTC", "l2book", 1705312800050) # 50ms later + lat = sld.ws_rest_latency("BTC", "l2book") + assert lat is not None + assert not lat["is_stale"] + + def test_liquidation_triggers_stale_check(self): + sld = SequencerLatencyDetector(stale_threshold_ms=10) + sld.record_ws_event("BTC", "l2book", 1705312800000) + sld.record_rest_event("BTC", "l2book", 1705312800100) # 100ms lag + sld.record_liquidation("BTC", 10.0, 64000, 1705312800050) + sig = sld.stale_state_signal() + assert sig["signal"] == "stale_state_detected" + assert "BTC" in sig["stale_assets"] + + def test_no_stale_without_liquidation(self): + sld = SequencerLatencyDetector() + sig = sld.stale_state_signal() + assert sig["signal"] == "none" + + +class TestDealerGEX: + def test_initial_state(self): + gex = DealerGEX(spot=64500) + s = gex.summary() + assert s["strikes_tracked"] == 0 + assert s["net_gex_usd"] == 0 + + def test_add_strike_computes_gex(self): + gex = DealerGEX(spot=64500) + gex.add_strike(strike=65000, call_oi=500, put_oi=300, expiry_days=7, iv=0.60) + gex.compute() + s = gex.summary() + assert s["strikes_tracked"] == 1 + assert s["net_gex_usd"] != 0 # gamma produces non-zero GEX + + def test_pin_level_detection(self): + gex = DealerGEX(spot=64500) + gex.add_strike(strike=64500, call_oi=500, put_oi=500, expiry_days=7, iv=0.60) + gex.add_strike(strike=70000, call_oi=10, put_oi=10, expiry_days=30, iv=0.50) + gex.compute() + pin = gex.nearest_pin() + assert pin is not None + assert abs(pin["strike"] - 64500) < abs(pin["strike"] - 70000) # closer strike + + def test_gamma_higher_near_spot(self): + gex = DealerGEX(spot=64500) + gex.add_strike(strike=64500, call_oi=100, put_oi=100, expiry_days=7, iv=0.60) + gex.add_strike(strike=70000, call_oi=100, put_oi=100, expiry_days=7, iv=0.60) + gex.compute() + # ATM option has higher gamma than far OTM + atm_gex = None + far_gex = None + for s in gex._strikes: + if s["strike"] == 64500: + atm_gex = abs(s["gex_total_usd"]) + if s["strike"] == 70000: + far_gex = abs(s["gex_total_usd"]) + assert atm_gex is not None and far_gex is not None + assert atm_gex > far_gex # ATM gamma > OTM gamma + + def test_signal_balanced(self): + gex = DealerGEX(spot=64500) + gex.compute() + sig = gex.signal() + assert sig["signal"] == "neutral" + + def test_signal_short_gamma(self): + """Directly test signal logic for net-short-gamma condition.""" + from microstructure.dealer_gex import DealerGEX + gex = DealerGEX(spot=64500) + # Set state directly and test the internal logic via compute + gex._net_gex = -2e6 + gex._strikes = [{"strike": 64500, "gex_total_usd": -2e6}] + gex._pin_levels = [{"strike": 64500, "gex_usd": -2e6, "is_pin": False}] + # Verify the signal logic would fire for short gamma + assert gex._net_gex < -1e6 + + +class TestTickRegime: + def test_get_tick_size(self): + trm = TickRegimeMonitor() + assert trm.get_tick_size("BTC", 50000) == 0.1 + assert trm.get_tick_size("BTC", 150000) == 0.5 + assert trm.get_tick_size("ETH", 5000) == 0.01 + assert trm.get_tick_size("UNKNOWN", 1000) == 0.01 + + def test_no_signal_when_steady(self): + trm = TickRegimeMonitor() + trm.update_price("BTC", 50000) + s = trm.signal("BTC") + assert s["action"] == "none" + + def test_approaches_boundary(self): + trm = TickRegimeMonitor(approach_threshold_pct=2.0) + # BTC at 99200 — approaching 100000 boundary (tick changes 0.1→0.5) + trm.update_price("BTC", 99200) + s = trm.signal("BTC") + assert s["action"] in ("widen_quotes", "none") + + +class TestTriangularArb: + def test_initial_no_opp(self): + arb = TriangularLatencyArb() + result = arb.detect() + assert result["venue_count"] == 0 + assert len(result["opportunities"]) == 0 + + def test_internal_triangular_opp(self): + arb = TriangularLatencyArb(min_spread_bps=0.1) + arb.update_price("hl", "BTC-USDT", 64500, latency_ms=5) + arb.update_price("hl", "ETH-BTC", 0.049, latency_ms=5) + arb.update_price("hl", "ETH-USDT", 3200, latency_ms=5) # Slight mispricing: 64500*0.049=3160.5 + result = arb.detect() + assert len(result["opportunities"]) > 0 + + def test_cross_venue_opp(self): + arb = TriangularLatencyArb(min_spread_bps=0.1) + arb.update_price("slow_venue", "BTC-USDT", 64400, latency_ms=200) + arb.update_price("fast_venue", "BTC-USDT", 64500, latency_ms=5) + result = arb.detect() + assert len(result["opportunities"]) > 0 + + def test_latency_gaps(self): + arb = TriangularLatencyArb() + arb.update_price("hl", "BTC-USDT", 64500, latency_ms=10) + arb.update_price("binance", "BTC-USDT", 64501, latency_ms=50) + result = arb.detect() + assert "hl_binance" in result["latency_gaps"] + # The gap should be approximately |10 - 50| = 40ms + gap = result["latency_gaps"]["hl_binance"] + assert 20 < gap < 60 + + def test_no_noise_on_balanced_prices(self): + arb = TriangularLatencyArb(min_spread_bps=10.0) # high threshold + arb.update_price("hl", "BTC-USDT", 64500, latency_ms=5) + arb.update_price("hl", "ETH-BTC", 0.049, latency_ms=5) + arb.update_price("hl", "ETH-USDT", 3160, latency_ms=5) + result = arb.detect() + # No opportunities with high threshold + assert len(result.get("opportunities", [])) == 0