diff --git a/live/filters/__init__.py b/live/filters/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/live/filters/toxicity.py b/live/filters/toxicity.py new file mode 100644 index 0000000..0b38000 --- /dev/null +++ b/live/filters/toxicity.py @@ -0,0 +1,116 @@ +""" +Pre-trade toxicity filter — blocks quoting when market is adversarially informed. + +Uses VPIN and flow toxicity metrics from microstructure module to +determine whether to quote, reduce size, or go flat. +""" + +from __future__ import annotations + +from collections import deque +from typing import Optional + + +class ToxicityFilter: + """VPIN-based flow toxicity filter for market-making. + + Monitors trade flow and order book to detect informed trading. + Blocks or reduces quoting when toxicity crosses thresholds. + + Usage: + tf = ToxicityFilter() + tf.update_trade(buy_vol=0.1, sell_vol=0.05) + decision = tf.check() + if decision["allow_quoting"]: + size_mult = decision["size_multiplier"] + """ + + def __init__( + self, + vpin_threshold: float = 0.3, + vpin_alarm: float = 0.5, + vpin_window: int = 50, + volume_bucket_size: float = 1.0, + size_reduction_steps: int = 5, + ): + self._threshold = vpin_threshold + self._alarm = vpin_alarm + self._window = vpin_window + self._bucket_size = volume_bucket_size + self._steps = size_reduction_steps + + self._buy_vol: deque[float] = deque(maxlen=5000) + self._sell_vol: deque[float] = deque(maxlen=5000) + self._current_vpin: float = 0.0 + self._last_obi: float = 0.0 + self._trade_count: int = 0 + + def update_trade(self, buy_vol: float, sell_vol: float): + """Feed a classified trade volume observation.""" + self._buy_vol.append(buy_vol) + self._sell_vol.append(sell_vol) + self._trade_count += 1 + + def update_book_imbalance(self, obi: float): + """Feed current order-book imbalance.""" + self._last_obi = obi + + def compute(self): + """Recompute VPIN from accumulated volume data.""" + from microstructure.toxicity import compute_vpin + result = compute_vpin( + list(self._buy_vol), + list(self._sell_vol), + volume_bucket_size=self._bucket_size, + n_buckets=self._window, + ) + self._current_vpin = result.get("vpin_value", 0.0) + + def check(self) -> dict: + """Evaluate current toxicity state. + + Returns: + allow_quoting: whether to quote at all + size_multiplier: fraction of base size to quote (1.0 = full, 0.0 = none) + vpin: current VPIN + obi: current OBI + reason: explanation if blocked/reduced + """ + self.compute() + v = self._current_vpin + + if v >= self._alarm: + return { + "allow_quoting": False, + "size_multiplier": 0.0, + "vpin": round(v, 4), + "obi": round(self._last_obi, 4), + "reason": f"VPIN {v:.3f} >= alarm {self._alarm}", + } + + if v >= self._threshold: + reduction = (v - self._threshold) / (self._alarm - self._threshold) + size = max(0.0, 1.0 - reduction) + return { + "allow_quoting": size > 0.0, + "size_multiplier": round(size, 2), + "vpin": round(v, 4), + "obi": round(self._last_obi, 4), + "reason": f"VPIN {v:.3f} >= threshold {self._threshold}", + } + + return { + "allow_quoting": True, + "size_multiplier": 1.0, + "vpin": round(v, 4), + "obi": round(self._last_obi, 4), + "reason": "ok", + } + + @property + def vpin(self) -> float: + return self._current_vpin + + @property + def trade_count(self) -> int: + return self._trade_count diff --git a/live/makers/__init__.py b/live/makers/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/live/makers/hl_btc_eth.py b/live/makers/hl_btc_eth.py new file mode 100644 index 0000000..b87fb49 --- /dev/null +++ b/live/makers/hl_btc_eth.py @@ -0,0 +1,218 @@ +""" +Hyperliquid BTC/ETH maker strategy with tight inventory limits. + +Uses Avellaneda-Stoikov optimal control from sim/maker.py, +toxicity filter from live/filters/toxicity.py, and treasury +from live/treasury.py for position/risk management. + +Designed for Phase 4 controlled deployment: post-only quotes, +tight inventory caps, toxicity gating. +""" + +from __future__ import annotations + +import time +from typing import Optional + +from sim.maker import AvellanedaStoikovMaker, MakerConfig, Quote +from sim.queue import QueueModel +from live.filters.toxicity import ToxicityFilter +from live.treasury import Treasury + + +class HlMaker: + """Hyperliquid market maker for a single coin. + + Lifecycle per tick: + 1. observe(mid_price) — feed mid price for vol estimation + 2. update_flow(buy_vol, sell_vol, obi) — feed trade flow for toxicity + 3. quote() — get bid/ask quotes (or None if blocked) + 4. record_fill(side, size, price, fee) — after exchange confirms fill + + Usage: + maker = HlMaker("BTC", treasury=treasury, max_inventory=0.003) + maker.observe(50000.0) + maker.update_flow(buy_vol=0.1, sell_vol=0.05, obi=0.2) + quote = maker.quote() + if quote: + # place bid at quote.bid, ask at quote.ask on exchange + ... + """ + + def __init__( + self, + coin: str, + treasury: Treasury, + max_inventory: float = 0.003, + base_size: float = 0.0002, + gamma: float = 0.1, + k: float = 1.5, + tau_hours: float = 1.0, + min_spread_bps: float = 1.0, + max_spread_bps: float = 15.0, + vpin_threshold: float = 0.3, + vpin_alarm: float = 0.5, + skew_factor: float = 0.3, + ): + self.coin = coin.upper() + self._treasury = treasury + self._max_inventory = max_inventory + self._base_size = base_size + self._skew_factor = skew_factor + + self._maker = AvellanedaStoikovMaker( + MakerConfig( + gamma=gamma, + k=k, + tau=tau_hours, + min_spread_bps=min_spread_bps, + max_spread_bps=max_spread_bps, + base_size=base_size, + max_inventory=max_inventory, + skew_factor=skew_factor, + ) + ) + + self._toxicity = ToxicityFilter( + vpin_threshold=vpin_threshold, + vpin_alarm=vpin_alarm, + ) + + self._mid_price: float = 0.0 + self._best_bid: float = 0.0 + self._best_ask: float = 0.0 + self._elapsed_hours: float = 0.0 + self._start_time: float = time.time() + self._last_quote: Optional[Quote] = None + + def observe(self, mid_price: float): + """Feed a new mid price observation.""" + self._mid_price = mid_price + self._elapsed_hours = (time.time() - self._start_time) / 3600.0 + self._maker.observe(mid_price) + self._treasury.update_mark_price(self.coin, mid_price) + + def update_book(self, best_bid: float, best_ask: float): + self._best_bid = best_bid + self._best_ask = best_ask + + def update_flow(self, buy_vol: float, sell_vol: float, obi: float = 0.0): + """Feed trade flow and order-book imbalance for toxicity tracking.""" + self._toxicity.update_trade(buy_vol, sell_vol) + self._toxicity.update_book_imbalance(obi) + + def quote(self) -> Optional[Quote]: + """Generate the next set of quotes, or None if blocked.""" + if self._treasury.is_halted(): + return None + if self._mid_price <= 0: + return None + + tox = self._toxicity.check() + if not tox["allow_quoting"]: + return None + + size_mult = tox["size_multiplier"] + + position = self._treasury.position(self.coin) + target_inv = 0.0 # neutral target + + q = self._maker.quote_with_skew( + mid_price=self._mid_price, + inventory=position, + elapsed_hours=self._elapsed_hours, + target_inventory=target_inv, + ) + + # Scale sizes by toxicity multiplier + q.bid_size *= size_mult + q.ask_size *= size_mult + + # Never cross the market + if self._best_bid > 0: + q.bid = round(min(q.bid, self._best_bid * 0.999), 2) + if self._best_ask > 0: + q.ask = round(max(q.ask, self._best_ask * 1.001), 2) + + self._last_quote = q + return q + + def record_fill(self, side: str, size: float, price: float, fee: float): + """Record a fill after exchange confirmation.""" + pnl = 0.0 + position = self._treasury.position(self.coin) + if (side == "sell" and position > 0) or (side == "buy" and position < 0): + pnl = size * (price - (self._mid_price)) + self._treasury.record_fill(self.coin, side, size, price, fee, pnl) + + def should_skip(self) -> bool: + """Check if we should skip quoting this tick.""" + if self._treasury.is_halted(): + return True + if self._treasury.position_size(self.coin) >= self._max_inventory: + return True + return False + + @property + def last_quote(self) -> Optional[Quote]: + return self._last_quote + + @property + def current_vpin(self) -> float: + return self._toxicity.vpin + + def summary(self) -> dict: + return { + "coin": self.coin, + "mid": self._mid_price, + "position": self._treasury.position(self.coin), + "vpin": self._toxicity.vpin, + "sigma": self._maker.sigma, + "last_quote": { + "bid": self._last_quote.bid, + "ask": self._last_quote.ask, + "spread_bps": self._last_quote.spread_bps, + } if self._last_quote else None, + } + + +class HlMakerPool: + """Manage multiple HlMaker instances across coins. + + Provides unified interface for multi-coin market making with + shared treasury and coordinated quoting. + """ + + def __init__(self, treasury: Treasury, maker_config: dict | None = None): + self._treasury = treasury + self._maker_config = maker_config or {} + self._makers: dict[str, HlMaker] = {} + + def add_maker(self, coin: str, max_inventory: float = 0.003, **kwargs) -> HlMaker: + cfg = dict(self._maker_config) + cfg.update(kwargs) + maker = HlMaker(coin=coin, treasury=self._treasury, max_inventory=max_inventory, **cfg) + self._makers[coin.upper()] = maker + return maker + + def get(self, coin: str) -> Optional[HlMaker]: + return self._makers.get(coin.upper()) + + def observe_all(self, mid_prices: dict[str, float]): + for coin, price in mid_prices.items(): + maker = self._makers.get(coin.upper()) + if maker: + maker.observe(price) + + def quote_all(self) -> dict[str, Optional[Quote]]: + return {coin: maker.quote() for coin, maker in self._makers.items()} + + def summary(self) -> dict: + return { + coin: maker.summary() + for coin, maker in self._makers.items() + } + + @property + def makers(self) -> dict[str, HlMaker]: + return self._makers diff --git a/live/monitors/__init__.py b/live/monitors/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/live/monitors/cross_venue.py b/live/monitors/cross_venue.py new file mode 100644 index 0000000..56c0133 --- /dev/null +++ b/live/monitors/cross_venue.py @@ -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 diff --git a/live/monitors/funding_basis.py b/live/monitors/funding_basis.py new file mode 100644 index 0000000..68f0ff0 --- /dev/null +++ b/live/monitors/funding_basis.py @@ -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} diff --git a/live/monitors/liq_risk.py b/live/monitors/liq_risk.py new file mode 100644 index 0000000..440167e --- /dev/null +++ b/live/monitors/liq_risk.py @@ -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 + ), + } diff --git a/live/treasury.py b/live/treasury.py new file mode 100644 index 0000000..3ee0c2e --- /dev/null +++ b/live/treasury.py @@ -0,0 +1,255 @@ +""" +Central treasury — position/capital limits, circuit breakers, PnL stops. + +Single source of truth for all risk constraints in live trading. +Integrates with sim/constraints.py for the constraint logic and adds +live-specific bookkeeping. +""" + +from __future__ import annotations + +import time +from typing import Optional + +from sim.constraints import ( + InventoryConstraint, + FundingConstraint, + FeeSchedule, + LiquidationRisk, + CircuitBreaker, +) + + +class Treasury: + """Central risk and capital management for live trading. + + Tracks: + - Current positions per asset + - Realized and unrealized PnL + - Daily trade counts + - Circuit breaker state + - Fee budget consumption + + Usage: + treasury = Treasury(initial_equity=10000.0) + ok = treasury.can_open("BTC", side="buy", size=0.001, mark_price=50000.0) + treasury.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0) + """ + + def __init__( + self, + initial_equity: float = 10000.0, + max_position_per_asset: float = 0.005, + max_net_exposure: float = 0.01, + max_daily_trades: int = 500, + max_drawdown_pct: float = -10.0, + max_toxic_rate: float = 0.4, + cooldown_seconds: float = 300.0, + maker_fee_pct: float = 0.0002, + taker_fee_pct: float = 0.0005, + maintenance_margin_pct: float = 0.03, + ): + self._initial_equity = initial_equity + self._realized_pnl: float = 0.0 + self._fees_paid: float = 0.0 + self._daily_trades: int = 0 + self._toxic_fills: int = 0 + self._api_errors: int = 0 + + # Positions tracked as {coin: {"side": "long"|"short", "size": float, "entry_px": float}} + self._positions: dict[str, dict] = {} + self._mark_prices: dict[str, float] = {} + + self._circuit_breaker = CircuitBreaker( + max_drawdown_pct=max_drawdown_pct, + max_daily_trades=max_daily_trades, + max_toxic_rate=max_toxic_rate, + cooldown_seconds=cooldown_seconds, + ) + + self._inventory = InventoryConstraint( + max_long=max_position_per_asset, + max_short=max_position_per_asset, + max_net_exposure=max_net_exposure, + ) + + self._fees = FeeSchedule(maker_fee_pct=maker_fee_pct, taker_fee_pct=taker_fee_pct) + self._funding = FundingConstraint() + self._liquidation = LiquidationRisk(maintenance_margin_pct=maintenance_margin_pct) + + self._halted: bool = False + self._halt_reason: str = "" + self._halted_at: float = 0.0 + self._session_start: float = time.time() + + # ── Position management ────────────────────────────────── + + def can_open(self, coin: str, side: str, size: float, mark_price: float) -> dict: + """Check whether a new position can be opened. + + Returns {allowed: bool, reason: str, fee_estimate: float} + """ + if self._halted: + return {"allowed": False, "reason": self._halt_reason, "fee_estimate": 0.0} + + pos = self._positions.get(coin.upper(), {}) + current_size = pos.get("size", 0.0) if pos.get("side") == side else -(pos.get("size", 0.0)) + new_size = current_size + size + + limits = self._inventory.check( + max(0.0, new_size) if side == "buy" else max(0.0, current_size), + max(0.0, -new_size) if side == "sell" else max(0.0, -current_size), + ) + + if not limits["long_ok"]: + return {"allowed": False, "reason": "long limit exceeded", "fee_estimate": 0.0} + if not limits["short_ok"]: + return {"allowed": False, "reason": "short limit exceeded", "fee_estimate": 0.0} + + fee = self._fees.maker_fee(size * mark_price) + return {"allowed": True, "reason": "ok", "fee_estimate": round(fee, 6)} + + def record_fill(self, coin: str, side: str, size: float, price: float, fee: float, pnl: float = 0.0): + """Record a filled trade.""" + c = coin.upper() + pos = self._positions.get(c) + + is_close = pos and pos.get("side") != side + if is_close: + self._realized_pnl += pnl + pos["size"] -= size + if pos["size"] <= 1e-10: + del self._positions[c] + else: + if not pos: + self._positions[c] = {"side": side, "size": size, "entry_px": price} + else: + total = pos["size"] + size + pos["entry_px"] = (pos["entry_px"] * pos["size"] + price * size) / total if total > 0 else price + pos["size"] = total + + self._fees_paid += fee + self._daily_trades += 1 + + self._check_breakers() + + def record_toxic_fill(self): + self._toxic_fills += 1 + + def record_api_error(self): + self._api_errors += 1 + + def update_mark_price(self, coin: str, price: float): + self._mark_prices[coin.upper()] = price + + # ── Position queries ───────────────────────────────────── + + def position(self, coin: str) -> float: + """Signed position (positive = long).""" + pos = self._positions.get(coin.upper(), {}) + raw = pos.get("size", 0.0) + return raw if pos.get("side") == "buy" else -raw + + def position_size(self, coin: str) -> float: + """Absolute position size.""" + return abs(self.position(coin)) + + @property + def all_positions(self) -> dict[str, float]: + return {c: self.position(c) for c in self._positions} + + @property + def net_exposure(self) -> float: + return sum(abs(p) for p in self.all_positions.values()) + + # ── PnL ────────────────────────────────────────────────── + + def unrealized_pnl(self) -> float: + pnl = 0.0 + for coin, pos in self._positions.items(): + mark = self._mark_prices.get(coin, pos.get("entry_px", 0)) + if pos["side"] == "buy": + pnl += pos["size"] * (mark - pos["entry_px"]) + else: + pnl += pos["size"] * (pos["entry_px"] - mark) + return round(pnl, 4) + + def total_pnl(self) -> float: + return self._realized_pnl + self.unrealized_pnl() - self._fees_paid + + def pnl_pct(self) -> float: + return self.total_pnl() / self._initial_equity * 100 if self._initial_equity > 0 else 0 + + @property + def equity(self) -> float: + return self._initial_equity + self.total_pnl() + + # ── Liquidation risk ───────────────────────────────────── + + def liquidation_distance(self, coin: str) -> float: + """Percentage distance to liquidation.""" + pos = self._positions.get(coin.upper()) + if not pos: + return float("inf") + + mark = self._mark_prices.get(coin.upper(), pos["entry_px"]) + liq = self._liquidation.liquidation_price( + entry_price=pos["entry_px"], + size=pos["size"], + position_side=pos["side"], + wallet_balance=self.equity, + ) + return self._liquidation.distance_to_liquidation_pct(mark, liq, pos["side"]) + + def is_liquidation_safe(self, coin: str, threshold_pct: float = 5.0) -> bool: + return self.liquidation_distance(coin) >= threshold_pct + + # ── Circuit breaker ────────────────────────────────────── + + def _check_breakers(self): + if self._halted: + return + toxic_rate = self._toxic_fills / max(self._daily_trades, 1) + result = self._circuit_breaker.evaluate({ + "pnl_pct": round(self.pnl_pct(), 2), + "daily_trades": self._daily_trades, + "toxic_rate": toxic_rate, + "api_errors": self._api_errors, + }) + if result.get("tripped"): + self._halted = True + self._halt_reason = result.get("reason", "unknown") + self._halted_at = time.time() + + def is_halted(self) -> bool: + if self._halted: + elapsed = time.time() - self._halted_at + if elapsed > self._circuit_breaker.cooldown_seconds: + self._halted = False + self._halt_reason = "" + self._daily_trades = 0 + self._toxic_fills = 0 + return self._halted + + @property + def halt_reason(self) -> str: + return self._halt_reason + + # ── Stats ──────────────────────────────────────────────── + + def summary(self) -> dict: + return { + "equity": round(self.equity, 2), + "realized_pnl": round(self._realized_pnl, 4), + "unrealized_pnl": self.unrealized_pnl(), + "total_pnl": self.total_pnl(), + "pnl_pct": round(self.pnl_pct(), 2), + "fees_paid": round(self._fees_paid, 4), + "daily_trades": self._daily_trades, + "toxic_fills": self._toxic_fills, + "api_errors": self._api_errors, + "positions": {c: round(v, 6) for c, v in self.all_positions.items()}, + "net_exposure": round(self.net_exposure, 6), + "halted": self._halted, + "uptime_hours": round((time.time() - self._session_start) / 3600, 1), + } diff --git a/tests/test_live_filters.py b/tests/test_live_filters.py new file mode 100644 index 0000000..192efa8 --- /dev/null +++ b/tests/test_live_filters.py @@ -0,0 +1,38 @@ +""" +Tests for live/filters/toxicity.py — pre-trade toxicity filter. +""" +from live.filters.toxicity import ToxicityFilter + + +class TestToxicityFilter: + def test_initial_state_allows_quoting(self): + tf = ToxicityFilter() + result = tf.check() + assert result["allow_quoting"] + assert result["size_multiplier"] == 1.0 + + def test_balanced_flow_keeps_quoting(self): + tf = ToxicityFilter() + for _ in range(200): + tf.update_trade(buy_vol=1.0, sell_vol=1.0) + result = tf.check() + assert result["allow_quoting"] + + def test_imbalanced_flow_reduces_or_blocks(self): + tf = ToxicityFilter(vpin_threshold=0.2, vpin_alarm=0.3, volume_bucket_size=1.0, vpin_window=10) + for _ in range(200): + tf.update_trade(buy_vol=3.0, sell_vol=1.0) + result = tf.check() + assert result["vpin"] >= 0 + + def test_book_imbalance_tracking(self): + tf = ToxicityFilter() + tf.update_book_imbalance(0.5) + result = tf.check() + assert result["obi"] == 0.5 + + def test_trade_count_increments(self): + tf = ToxicityFilter() + tf.update_trade(0.1, 0.05) + tf.update_trade(0.2, 0.1) + assert tf.trade_count == 2 diff --git a/tests/test_live_maker.py b/tests/test_live_maker.py new file mode 100644 index 0000000..c408487 --- /dev/null +++ b/tests/test_live_maker.py @@ -0,0 +1,86 @@ +""" +Tests for live/makers/hl_btc_eth.py — HL maker strategy. +""" +from live.treasury import Treasury +from live.makers.hl_btc_eth import HlMaker, HlMakerPool + + +class TestHlMaker: + def test_initial_quote_returns_none(self): + t = Treasury(initial_equity=10000.0) + maker = HlMaker("BTC", treasury=t) + assert maker.quote() is None # no mid price yet + + def test_quote_after_observe(self): + t = Treasury(initial_equity=10000.0) + maker = HlMaker("BTC", treasury=t, base_size=0.001) + maker.observe(50000.0) + maker.update_book(49999.0, 50001.0) + q = maker.quote() + assert q is not None + assert q.bid < 50000.0 < q.ask + assert q.bid_size > 0 + assert q.ask_size > 0 + + def test_quote_blocked_by_halte(self): + t = Treasury(initial_equity=10000.0, max_drawdown_pct=-1.0) + t._realized_pnl = -5000.0 + t._fees_paid = 0 + t._check_breakers() + maker = HlMaker("BTC", treasury=t) + maker.observe(50000.0) + assert maker.quote() is None + + def test_should_skip_at_max_inventory(self): + t = Treasury(initial_equity=10000.0) + maker = HlMaker("BTC", treasury=t, max_inventory=0.001) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + maker.observe(50000.0) + assert maker.should_skip() + + def test_record_fill_updates_treasury(self): + t = Treasury(initial_equity=10000.0) + maker = HlMaker("BTC", treasury=t) + maker.observe(50000.0) + maker.record_fill("buy", 0.001, 50000.0, 10.0) + assert t.position("BTC") == 0.001 + + def test_quote_never_crosses_book(self): + t = Treasury(initial_equity=10000.0) + maker = HlMaker("BTC", treasury=t, base_size=0.001, gamma=0.5) + maker.observe(50000.0) + maker.update_book(49995.0, 50005.0) + for _ in range(20): + q = maker.quote() + if q: + assert q.bid <= 49995.0 # never above best bid + assert q.ask >= 50005.0 # never below best ask + assert q.bid < q.ask + + +class TestHlMakerPool: + def test_add_and_get_makers(self): + t = Treasury(initial_equity=10000.0) + pool = HlMakerPool(treasury=t) + pool.add_maker("BTC", max_inventory=0.002) + pool.add_maker("ETH", max_inventory=0.01) + assert pool.get("BTC") is not None + assert pool.get("ETH") is not None + assert pool.get("SOL") is None + + def test_observe_and_quote_all(self): + t = Treasury(initial_equity=10000.0) + pool = HlMakerPool(treasury=t) + maker_btc = pool.add_maker("BTC", max_inventory=0.002, base_size=0.001) + pool.observe_all({"BTC": 50000.0}) + quotes = pool.quote_all() + assert "BTC" in quotes + assert quotes["BTC"] is not None + + def test_summary(self): + t = Treasury(initial_equity=10000.0) + pool = HlMakerPool(treasury=t) + pool.add_maker("BTC", max_inventory=0.002) + pool.observe_all({"BTC": 50000.0}) + s = pool.summary() + assert "BTC" in s diff --git a/tests/test_live_monitors.py b/tests/test_live_monitors.py new file mode 100644 index 0000000..b9ac327 --- /dev/null +++ b/tests/test_live_monitors.py @@ -0,0 +1,93 @@ +""" +Tests for live/monitors — cross-venue, funding/basis, liquidation risk. +""" +from live.monitors.cross_venue import CrossVenueMonitor +from live.monitors.funding_basis import FundingBasisMonitor +from live.monitors.liq_risk import LiquidationRiskOverlay +from live.treasury import Treasury + + +class TestCrossVenueMonitor: + def test_update_and_lead_lag(self): + cm = CrossVenueMonitor(window=50, max_lag=5) + for i in range(50): + cm.update("hl", "BTC", 50000.0 + i * 10, float(i)) + cm.update("binance", "BTC", 50000.0 + i * 10 + 2, float(i)) + result = cm.lead_lag("BTC", "hl", "binance") + assert result is not None + assert "correlation" in result + assert "lag" in result + + def test_spot_premium(self): + cm = CrossVenueMonitor() + for _ in range(10): + cm.update("hl", "BTC", 50005.0, 0) + cm.update("binance", "BTC", 50000.0, 0) + premium = cm.spot_premium("BTC") + assert premium is not None + assert premium["basis_bps"] > 0 + + def test_nonexistent_coin_returns_none(self): + cm = CrossVenueMonitor() + assert cm.lead_lag("XYZ", "hl", "binance") is None + + def test_summary(self): + cm = CrossVenueMonitor() + for i in range(50): + cm.update("hl", "BTC", 50000.0 + i * 10, float(i)) + cm.update("binance", "BTC", 50000.0 + i * 10, float(i)) + s = cm.summary("BTC") + assert "hl_binance" in s + + +class TestFundingBasisMonitor: + def test_initial_no_signal(self): + fm = FundingBasisMonitor() + result = fm.signal("BTC") + assert result["signal"] == "insufficient_data" + + def test_signal_with_data(self): + fm = FundingBasisMonitor(funding_window=100, samples_per_hour=60) + for _ in range(200): + fm.update_funding("BTC", 0.00001) + result = fm.signal("BTC") + assert result["signal"] in ("neutral", "positive", "negative", "high_positive", "high_negative") + assert "funding_mean_annual_pct" in result + + def test_basis_requires_spot_and_perp(self): + fm = FundingBasisMonitor() + for i in range(50): + fm.update_perp("BTC", 50005.0) + fm.update_spot("BTC", 50000.0) + fm.update_funding("BTC", 0.00001) + result = fm.signal("BTC") + assert result["basis_current_bps"] > 0 + + +class TestLiquidationRiskOverlay: + def test_safe_position(self): + t = Treasury(initial_equity=100000.0) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + overlay = LiquidationRiskOverlay(treasury=t) + result = overlay.check("BTC") + assert result["level"] == "safe" + + def test_no_position(self): + t = Treasury(initial_equity=10000.0) + overlay = LiquidationRiskOverlay(treasury=t) + result = overlay.check("BTC") + assert result["distance_pct"] > 1e5 # capped at 999999 for display + + def test_recommended_action(self): + t = Treasury(initial_equity=10000.0) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + overlay = LiquidationRiskOverlay(treasury=t) + assert overlay.recommended_action("BTC") == "none" + + def test_summary(self): + t = Treasury(initial_equity=100000.0) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + overlay = LiquidationRiskOverlay(treasury=t) + s = overlay.summary() + assert "positions" in s + assert "worst_case" in s diff --git a/tests/test_live_treasury.py b/tests/test_live_treasury.py new file mode 100644 index 0000000..40a3f67 --- /dev/null +++ b/tests/test_live_treasury.py @@ -0,0 +1,95 @@ +""" +Tests for live/treasury.py — central treasury and risk management. +""" +from live.treasury import Treasury + + +class TestTreasury: + def test_initial_equity(self): + t = Treasury(initial_equity=10000.0) + assert t.equity == 10000.0 + assert t.total_pnl() == 0.0 + + def test_can_open_within_limits(self): + t = Treasury(initial_equity=10000.0, max_position_per_asset=0.005) + result = t.can_open("BTC", side="buy", size=0.001, mark_price=50000.0) + assert result["allowed"] + + def test_cannot_open_exceed_inventory(self): + t = Treasury(initial_equity=10000.0, max_position_per_asset=0.002) + t.can_open("BTC", side="buy", size=0.001, mark_price=50000.0) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + result = t.can_open("BTC", side="buy", size=0.0015, mark_price=50000.0) + assert not result["allowed"] + + def test_record_fill_updates_position(self): + t = Treasury(initial_equity=10000.0) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + assert t.position("BTC") == 0.001 + + def test_record_close_updates_pnl(self): + t = Treasury(initial_equity=10000.0) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + t.record_fill("BTC", side="sell", size=0.001, price=50100.0, fee=10.0, pnl=100.0) + assert t.position("BTC") == 0.0 + assert t.total_pnl() > 0 + + def test_unrealized_pnl(self): + t = Treasury(initial_equity=10000.0) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + t.update_mark_price("BTC", 50200.0) + assert t.unrealized_pnl() > 0 + + def test_pnl_pct(self): + t = Treasury(initial_equity=10000.0) + t.record_fill("BTC", side="buy", size=0.1, price=50000.0, fee=0.1, pnl=0) + t.update_mark_price("BTC", 50200.0) # 0.1 * 200 = $20 unrealized >> $0.10 fee + assert t.pnl_pct() > 0 + + def test_circuit_breaker_drawdown(self): + t = Treasury(initial_equity=10000.0, max_drawdown_pct=-1.0) + # Force large negative PnL + t._realized_pnl = -5000.0 + t._fees_paid = 0 + t._check_breakers() + assert t.is_halted() + + def test_liquidation_distance(self): + t = Treasury(initial_equity=50000.0) + t.record_fill("BTC", side="buy", size=1.0, price=50000.0, fee=10.0, pnl=0) + dist = t.liquidation_distance("BTC") + assert dist > 0 + + def test_all_positions(self): + t = Treasury() + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + t.record_fill("ETH", side="sell", size=0.01, price=3000.0, fee=10.0, pnl=0) + positions = t.all_positions + assert positions["BTC"] == 0.001 + assert positions["ETH"] == -0.01 + + def test_net_exposure(self): + t = Treasury() + t.record_fill("BTC", side="buy", size=0.002, price=50000.0, fee=10.0, pnl=0) + t.record_fill("ETH", side="buy", size=0.003, price=3000.0, fee=10.0, pnl=0) + assert t.net_exposure == 0.005 + + def test_summary(self): + t = Treasury(initial_equity=10000.0) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + s = t.summary() + assert "equity" in s + assert "pnl_pct" in s + assert "positions" in s + assert "BTC" in s["positions"] + + def test_toxic_fill_tracking(self): + t = Treasury(initial_equity=10000.0, max_toxic_rate=0.1) + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=0) + for _ in range(9): + t.record_fill("BTC", side="buy", size=0.001, price=50000.0, fee=10.0, pnl=-1.0) + t.record_toxic_fill() + t.record_toxic_fill() + t._check_breakers() + # 2/10 = 20% toxic > 10% threshold + assert t.is_halted()