feat: Phase 4 — controlled strategy deployment module + 38 tests

New live/ sub-modules for production-ready market making:

live/filters/toxicity.py (ToxicityFilter):
  VPIN-based pre-trade filter. Accumulates buy/sell volume, computes
  VPIN via microstructure module, produces quoting decision:
    - allow_quoting: bool
    - size_multiplier: 0.0–1.0 (graduated reduction approaching alarm)
    - granular thresholds (threshold vs alarm) with smooth reduction

live/treasury.py (Treasury):
  Central capital/risk management — single source of truth:
  - Position tracking per coin (opening, closing, average entry)
  - Realized + unrealized PnL computation
  - Pre-trade constraint checks (inventory limits, fee estimates)
  - Circuit breaker (drawdown, trade count, toxic fill rate, API errors)
  - Liquidation distance monitoring
  - Automatic cooldown reset after trip expiry

live/makers/hl_btc_eth.py:
  HlMaker — per-coin market maker integrating:
    - AvellanedaStoikovMaker (Phase 3) for optimal quotes
    - ToxicityFilter for pre-trade gating
    - Treasury for position/risk checks
  HlMakerPool — manages multiple HlMaker instances with shared treasury
    and coordinated observe_all()/quote_all()

live/monitors/cross_venue.py (CrossVenueMonitor):
  Cross-exchange lead-lag detection via cross-correlation at multiple
  lags. Spot premium (basis proxy) computation. Multi-venue summary.

live/monitors/funding_basis.py (FundingBasisMonitor):
  Funding regime classification, momentum detection, carry PnL
  estimation, basis spread analysis. Uses microstructure/funding.py.

live/monitors/liq_risk.py (LiquidationRiskOverlay):
  Per-position liquidation distance monitoring with tiered warnings
  (safe/warning/danger/critical). Recommended position reduction.

38 tests across 4 files (all pass):
  test_live_filters.py (5)
  test_live_maker.py (9)
  test_live_monitors.py (11)
  test_live_treasury.py (13)

Total test suite: 172 tests, all passing.
This commit is contained in:
ramseshk
2026-08-07 14:47:08 +08:00
parent 639dd4fb6d
commit 4f66ef36a9
13 changed files with 1229 additions and 0 deletions
View File
+120
View File
@@ -0,0 +1,120 @@
"""
Cross-venue lead-lag monitor.
Detects when one exchange leads another in price discovery.
Used for informational purposes only in Phase 4 — no auto-trading.
"""
from __future__ import annotations
from collections import deque
from typing import Optional
import numpy as np
class CrossVenueMonitor:
"""Monitor price lead-lag relationships between exchanges.
Tracks mid prices across venues and computes cross-correlation
and lead-lag structure. Can detect when Hyperliquid follows
Binance or vice versa.
Usage:
monitor = CrossVenueMonitor(pairs=[("hl", "binance")], window=100)
monitor.update("hl", "BTC", 50000.0)
monitor.update("binance", "BTC", 50000.5)
result = monitor.lead_lag("BTC")
"""
def __init__(
self,
pairs: list[tuple[str, str]] | None = None,
window: int = 100,
max_lag: int = 10,
):
self._window = window
self._max_lag = max_lag
self._pairs = pairs or [("hl", "binance"), ("hl", "bybit"), ("hl", "okx")]
# venue → coin → deque of mid prices
self._prices: dict[str, dict[str, deque[float]]] = {}
self._timestamps: dict[str, dict[str, deque[float]]] = {}
def update(self, venue: str, coin: str, price: float, timestamp: float):
"""Record a mid price observation from a venue."""
v = venue.lower()
c = coin.upper()
self._prices.setdefault(v, {}).setdefault(c, deque(maxlen=self._window)).append(price)
self._timestamps.setdefault(v, {}).setdefault(c, deque(maxlen=self._window)).append(timestamp)
def lead_lag(self, coin: str, venue_a: str = "hl", venue_b: str = "binance") -> dict | None:
"""Determine which venue leads by cross-correlation at various lags.
Returns {leading_venue: str, max_correlation: float, lag: int}
Negative lag = venue_a leads, positive lag = venue_b leads.
"""
prices_a = list(self._prices.get(venue_a.lower(), {}).get(coin.upper(), []))
prices_b = list(self._prices.get(venue_b.lower(), {}).get(coin.upper(), []))
min_len = min(len(prices_a), len(prices_b))
if min_len < self._max_lag + 2:
return None
a = np.array(prices_a[-min_len:])
b = np.array(prices_b[-min_len:])
best_corr = -1.0
best_lag = 0
for lag in range(-self._max_lag, self._max_lag + 1):
if lag < 0:
corr = np.corrcoef(a[-lag:], b[:lag])[0, 1] if lag < 0 else 0
elif lag > 0:
corr = np.corrcoef(a[:min_len - lag], b[lag:])[0, 1]
else:
corr = np.corrcoef(a, b)[0, 1]
if not np.isnan(corr) and abs(corr) > abs(best_corr):
best_corr = float(corr)
best_lag = lag
return {
"leading_venue": venue_a if best_lag < 0 else venue_b if best_lag > 0 else "none",
"correlation": round(best_corr, 4),
"lag": best_lag,
"samples": min_len,
}
def spot_premium(self, coin: str, venue: str = "hl", spot_venue: str = "binance") -> dict | None:
"""Compute the premium of venue over spot (basis proxy)."""
v = self._prices.get(venue.lower(), {}).get(coin.upper())
sv = self._prices.get(spot_venue.lower(), {}).get(coin.upper())
if not v or not sv:
return None
perp = v[-1]
spot = sv[-1]
basis_bps = (perp - spot) / spot * 10000 if spot > 0 else 0
return {
"venue": venue,
"spot_venue": spot_venue,
"perp_price": perp,
"spot_price": spot,
"basis_bps": round(basis_bps, 2),
}
def summary(self, coin: str) -> dict:
"""Summary for a given coin across all venues."""
result = {}
for va, vb in self._pairs:
ll = self.lead_lag(coin, va, vb)
if ll:
result[f"{va}_{vb}"] = ll
premium = self.spot_premium(coin)
if premium:
result["premium"] = premium
return result
+99
View File
@@ -0,0 +1,99 @@
"""
Funding and basis carry monitor.
Tracks funding rates across Hyperliquid and estimates carry
trade profitability. Signals when funding arbitrage is attractive.
Phase 4: informational only — no automated trading.
"""
from __future__ import annotations
from collections import deque
from typing import Optional
from microstructure.funding import (
funding_regime,
funding_predictability,
funding_carry_pnl,
)
class FundingBasisMonitor:
"""Monitor funding rates and carry trade opportunities.
Usage:
monitor = FundingBasisMonitor()
monitor.update_funding("BTC", 0.0001)
monitor.update_spot("BTC", 50000.0)
monitor.update_perp("BTC", 50005.0)
result = monitor.signal("BTC")
"""
def __init__(
self,
funding_window: int = 1440, # 24h at 1 sample/min
samples_per_hour: int = 60,
):
self._samples_per_hour = samples_per_hour
self._funding: dict[str, deque[float]] = {}
self._perp_prices: dict[str, deque[float]] = {}
self._spot_prices: dict[str, deque[float]] = {}
self._window = funding_window
def update_funding(self, coin: str, rate: float):
"""Record an hourly funding rate."""
self._funding.setdefault(coin.upper(), deque(maxlen=self._window)).append(rate)
def update_perp(self, coin: str, price: float):
self._perp_prices.setdefault(coin.upper(), deque(maxlen=self._window)).append(price)
def update_spot(self, coin: str, price: float):
self._spot_prices.setdefault(coin.upper(), deque(maxlen=self._window)).append(price)
def signal(self, coin: str) -> dict:
"""Generate funding/basis signal for a coin."""
c = coin.upper()
funding_list = list(self._funding.get(c, []))
perp_list = list(self._perp_prices.get(c, []))
spot_list = list(self._spot_prices.get(c, []))
if not funding_list:
return {"signal": "insufficient_data", "action": "none"}
regime = funding_regime(funding_list, window_hours=min(24, len(funding_list) // self._samples_per_hour),
n_samples_per_hour=self._samples_per_hour)
predictability = funding_predictability(funding_list)
basis = None
if perp_list and spot_list:
from microstructure.funding import basis_spread
basis = basis_spread(perp_list, spot_list)
carry = funding_carry_pnl(
funding_list,
position_size=1.0,
mark_prices=perp_list if perp_list else None,
n_samples_per_hour=self._samples_per_hour,
)
regime_name = regime.get("regime", "unknown")
action = "none"
if regime_name in ("high_positive",) and predictability.get("is_momentum"):
action = "consider_short" # shorts earn positive funding
elif regime_name in ("high_negative",) and predictability.get("is_momentum"):
action = "consider_long"
return {
"signal": regime_name,
"action": action,
"funding_mean_annual_pct": regime.get("mean_annual_pct", 0),
"momentum": predictability.get("momentum_strength", 0),
"carry_cumulative_pnl": carry.get("cumulative_pnl", 0),
"basis_current_bps": basis.get("current_basis_bps", 0) if basis else 0,
"basis_mean_bps": basis.get("mean_basis_bps", 0) if basis else 0,
}
def summary(self) -> dict:
return {coin: self.signal(coin) for coin in self._funding}
+109
View File
@@ -0,0 +1,109 @@
"""
Liquidation risk overlay.
Monitors current positions, mark prices, and computes liquidation
distance. Warns when positions approach liquidation threshold.
Integrates with live/treasury.py for position tracking.
"""
from __future__ import annotations
from typing import Optional
from live.treasury import Treasury
class LiquidationRiskOverlay:
"""Liquidation risk monitor for open positions.
Usage:
overlay = LiquidationRiskOverlay(treasury=treasury)
status = overlay.check("BTC")
if status["warning"]:
# reduce position or add margin
"""
def __init__(
self,
treasury: Treasury,
warning_threshold_pct: float = 10.0,
danger_threshold_pct: float = 5.0,
critical_threshold_pct: float = 2.5,
):
self._treasury = treasury
self._warning = warning_threshold_pct
self._danger = danger_threshold_pct
self._critical = critical_threshold_pct
def check(self, coin: str) -> dict:
"""Check liquidation safety for a specific coin."""
distance = self._treasury.liquidation_distance(coin.upper())
if distance >= self._warning or distance >= 1e9:
level = "safe"
elif distance >= self._danger:
level = "warning"
elif distance >= self._critical:
level = "danger"
else:
level = "critical"
return {
"coin": coin.upper(),
"level": level,
"distance_pct": round(min(distance, 999999), 2),
"position": self._treasury.position(coin.upper()),
"warning": level in ("warning", "danger", "critical"),
"needs_action": level == "critical",
}
def check_all(self) -> dict[str, dict]:
positions = self._treasury.all_positions
return {coin: self.check(coin) for coin, pos in positions.items() if abs(pos) > 0}
def pnl_at_liquidation(self, coin: str) -> float:
"""Estimate realized PnL if position reaches liquidation price."""
pos = self._treasury._positions.get(coin.upper())
if not pos:
return 0.0
entry = pos["entry_px"]
size = pos["size"]
side = pos["side"]
liq_price = self._treasury._liquidation.liquidation_price(
entry, size, side, self._treasury.equity
)
if side == "buy":
return size * (liq_price - entry)
else:
return size * (entry - liq_price)
def recommended_action(self, coin: str) -> str:
"""Recommend action based on liquidation distance."""
status = self.check(coin)
level = status["level"]
if level == "safe":
return "none"
elif level == "warning":
return "reduce_position_25pct"
elif level == "danger":
return "reduce_position_50pct"
else:
return "close_all"
def summary(self) -> dict:
return {
"positions": self.check_all(),
"worst_case": min(
(self.check(c)["distance_pct"] for c in self._treasury.all_positions),
default=float("inf")
),
"any_critical": any(
self.check(c)["level"] == "critical"
for c in self._treasury.all_positions
),
}