diff --git a/backtests/tick_runner.py b/backtests/tick_runner.py new file mode 100644 index 0000000..1788911 --- /dev/null +++ b/backtests/tick_runner.py @@ -0,0 +1,519 @@ +""" +Tick-level backtest runner — replays real Parquet L2/trade events +through the event-driven simulation engine. + +This is the foundation for HFT strategy validation. Unlike the VBT runner +which uses candle data, this replays every L2 snapshot, trade tick, and +mark price update sequentially through the queue-position-aware simulator. + +Usage: + python backtests/tick_runner.py --coin BTC --start 2026-08-01 --end 2026-08-07 + python backtests/tick_runner.py --coin ETH --days 3 --maker as_mm --gamma 0.15 + python backtests/tick_runner.py --coin BTC --maker vpin_as_mm --vpin-threshold 0.3 +""" + +from __future__ import annotations + +import argparse +import json +import logging +import math +import os +import sys +from datetime import datetime, timezone +from pathlib import Path +from typing import Optional + +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) + +from sim.engine import SimulationEngine, SimConfig +from sim.maker import AvellanedaStoikovMaker, MakerConfig +from sim.fills import FillModelConfig +from sim.scenario import ScenarioConfig +from quant.significance import QuantVerdict + +logger = logging.getLogger(__name__) + +RESULTS_DIR = Path(__file__).resolve().parent / "results" / "tick" +RESULTS_DIR.mkdir(parents=True, exist_ok=True) + + +class VPINGatedASMaker(AvellanedaStoikovMaker): + """A-S market maker with VPIN toxicity gating and inventory skew. + + Inherits the base A-S quoting logic and adds: + - VPIN gating: stop quoting when VPIN exceeds threshold + - Inventory skew: bias quotes toward reducing inventory + - Dynamic spread: widen spread when VPIN is elevated but below alarm + """ + + def __init__( + self, + config: MakerConfig | None = None, + vpin_threshold: float = 0.30, + vpin_alarm: float = 0.50, + vpin_window: int = 50, + ): + super().__init__(config) + self._vpin_threshold = vpin_threshold + self._vpin_alarm = vpin_alarm + self._current_vpin: float = 0.0 + self._buy_vol: list[float] = [] + self._sell_vol: list[float] = [] + self._vpin_window = vpin_window + self._vpins: list[float] = [] + + def update_vpin(self, buy_vol: float, sell_vol: float): + self._buy_vol.append(buy_vol) + self._sell_vol.append(sell_vol) + if len(self._buy_vol) > self._vpin_window * 20: + self._buy_vol = self._buy_vol[-self._vpin_window * 20:] + self._sell_vol = self._sell_vol[-self._vpin_window * 20:] + self._recompute_vpin() + + def _recompute_vpin(self): + from microstructure.toxicity import compute_vpin + result = compute_vpin( + list(self._buy_vol), list(self._sell_vol), + n_buckets=self._vpin_window, + ) + self._current_vpin = result.get("vpin_value", 0.0) + self._vpins.append(self._current_vpin) + if len(self._vpins) > 200: + self._vpins = self._vpins[-200:] + + @property + def vpin(self) -> float: + return self._current_vpin + + def allowed_to_quote(self) -> tuple[bool, float]: + if self._current_vpin >= self._vpin_alarm: + return False, 0.0 + if self._current_vpin >= self._vpin_threshold: + reduction = (self._current_vpin - self._vpin_threshold) / ( + self._vpin_alarm - self._vpin_threshold + ) + return True, max(0.0, 1.0 - reduction) + return True, 1.0 + + def quote( + self, + mid_price: float, + inventory: float, + elapsed_hours: float, + ) -> Optional["Quote"]: + from sim.maker import Quote + + allowed, size_mult = self.allowed_to_quote() + if not allowed or size_mult <= 0: + return None + + base = super().quote(mid_price, inventory, elapsed_hours) + + inv_ratio = inventory / max(self._cfg.max_inventory, 0.0001) + skew = inv_ratio * self._cfg.skew_factor * base.spread_bps / 10000 * mid_price + + spread_mult = 1.0 + if self._current_vpin >= self._vpin_threshold * 0.7: + spread_mult = 1.0 + (self._current_vpin - self._vpin_threshold * 0.7) / ( + self._vpin_alarm - self._vpin_threshold * 0.7 + ) + + bid = base.bid - skew + ask = base.ask - skew + half_spread = base.spread_bps / 10000 * mid_price * spread_mult / 2 + bid = mid_price + (bid - mid_price) - half_spread * (spread_mult - 1) + ask = mid_price + (ask - mid_price) + half_spread * (spread_mult - 1) + + if inventory > 0: + bid_size = base.bid_size * size_mult * (1.0 - inv_ratio * self._cfg.skew_factor) + ask_size = base.ask_size * size_mult * (1.0 + inv_ratio * self._cfg.skew_factor) + else: + bid_size = base.bid_size * size_mult * (1.0 + abs(inv_ratio) * self._cfg.skew_factor) + ask_size = base.ask_size * size_mult * (1.0 - abs(inv_ratio) * self._cfg.skew_factor) + + bid = max(bid, 1.0) + ask = max(ask, bid + mid_price * self._cfg.min_spread_bps / 10000) + bid_size = max(bid_size, self._cfg.base_size * 0.1) + ask_size = max(ask_size, self._cfg.base_size * 0.1) + + new_spread_bps = (ask - bid) / mid_price * 10000 if mid_price > 0 else 0 + + return Quote( + bid=round(bid, 2), + ask=round(ask, 2), + bid_size=round(bid_size, 6), + ask_size=round(ask_size, 6), + reservation=round(base.reservation, 2), + spread_bps=round(new_spread_bps, 2), + ) + + +class TickBacktestRunner: + """Replays stored Parquet L2/trade events through the simulation engine. + + The engine processes events sequentially: + 1. L2 update → update book, maybe re-quote + 2. Trade → check fills, update PnL + 3. Periodic → funding tick, re-quote, cancel stale orders + """ + + def __init__( + self, + data_dir: str = "data/raw", + maker_type: str = "as_mm", + maker_config: MakerConfig | None = None, + vpin_threshold: float = 0.30, + vpin_alarm: float = 0.50, + gamma: float = 0.1, + base_size: float = 0.001, + max_inventory: float = 0.005, + skew_factor: float = 0.5, + maker_fee_pct: float = 0.0002, + taker_fee_pct: float = 0.0005, + adverse_selection_prob: float = 0.15, + cancel_after_ms: float = 5000.0, + quote_refresh_ms: float = 2000.0, + seed: int | None = 42, + ): + self._data_dir = data_dir + self._maker_type = maker_type + self._vpin_threshold = vpin_threshold + self._vpin_alarm = vpin_alarm + self._maker_fee_pct = maker_fee_pct + self._taker_fee_pct = taker_fee_pct + + maker_cfg = maker_config or MakerConfig( + gamma=gamma, + base_size=base_size, + max_inventory=max_inventory, + skew_factor=skew_factor, + ) + + if maker_type == "vpin_as_mm": + self._maker = VPINGatedASMaker( + config=maker_cfg, + vpin_threshold=vpin_threshold, + vpin_alarm=vpin_alarm, + ) + else: + self._maker = AvellanedaStoikovMaker(maker_cfg) + + self._sim_config = SimConfig( + maker=maker_cfg, + fills=FillModelConfig(adverse_selection_prob=adverse_selection_prob), + maker_fee_pct=maker_fee_pct, + taker_fee_pct=taker_fee_pct, + cancel_after_ms=cancel_after_ms, + quote_refresh_ms=quote_refresh_ms, + seed=seed, + ) + + self._engine: Optional[SimulationEngine] = None + self._seed = seed + + def load_events( + self, + coin: str, + start_date: str, + end_date: str, + ) -> list[dict]: + """Load L2 and trade events from Parquet store and merge into a single timeline. + + Returns a list of events sorted by timestamp, each with: + {"type": "l2"|"trade", "data": {...}, "time": float, "coin": str} + """ + from data.store import read_range + + logger.info("Loading L2 data for %s from %s to %s...", coin, start_date, end_date) + l2_msgs = read_range(self._data_dir, channel="l2book", coin=coin.upper(), + start_date=start_date, end_date=end_date) + logger.info("Loaded %d L2 messages", len(l2_msgs)) + + logger.info("Loading trade data for %s from %s to %s...", coin, start_date, end_date) + trade_msgs = read_range(self._data_dir, channel="trades", coin=coin.upper(), + start_date=start_date, end_date=end_date) + logger.info("Loaded %d trade messages", len(trade_msgs)) + + events = [] + + t0 = None + for msg in l2_msgs: + payload = msg.get("payload", {}) + local_ts = msg.get("local_ts", 0) + if not t0: + t0 = local_ts + + bids = {} + asks = {} + levels = payload.get("levels", []) + msg_type = payload.get("type", "snapshot") + + if msg_type == "snapshot" and isinstance(levels, list): + if len(levels) >= 1: + for bid in levels[0]: + sz = float(bid.get("sz", 0)) + if sz > 0: + bids[float(bid["px"])] = sz + if len(levels) >= 2: + for ask in levels[1]: + sz = float(ask.get("sz", 0)) + if sz > 0: + asks[float(ask["px"])] = sz + elif msg_type == "delta": + delta = payload.get("delta", {}) + if delta: + px = float(delta.get("px", 0)) + sz = float(delta.get("sz", 0)) + side = delta.get("side", "B") + if sz <= 0: + continue + bids = {px: sz} if side == "B" else {} + asks = {px: sz} if side == "A" else {} + + if bids or asks: + events.append({ + "type": "l2", + "data": {"bids": bids, "asks": asks}, + "time": local_ts - t0 if t0 else local_ts, + "coin": coin.upper(), + }) + + for msg in trade_msgs: + payload = msg.get("payload", {}) + local_ts = msg.get("local_ts", 0) + events.append({ + "type": "trade", + "data": { + "px": payload.get("px", "0"), + "sz": payload.get("sz", "0"), + "side": payload.get("side", "B"), + }, + "time": local_ts - t0 if t0 else local_ts, + "coin": coin.upper(), + }) + + events.sort(key=lambda e: e["time"]) + logger.info("Total events: %d (%.1f hours)", len(events), + (events[-1]["time"] - events[0]["time"]) / 3600 if events else 0) + return events + + def run( + self, + coin: str, + start_date: str, + end_date: str, + ) -> dict: + """Load events and run the simulation engine. Returns result dict.""" + events = self.load_events(coin, start_date, end_date) + if not events or len(events) < 2: + logger.error("Not enough events to run backtest") + return self._empty_result(coin, start_date, end_date) + + engine = SimulationEngine( + config=self._sim_config, + maker=self._maker, + seed=self._seed, + ) + + engine.run(events) + self._engine = engine + + stats = engine.stats() + breakdown = engine.breakdown() + + equity_curve = engine.reporter.equity_curve + total_trades = stats.total_trades + duration_hours = (events[-1]["time"] - events[0]["time"]) / 3600 if events else 0 + + returns = [] + eq_vals = [p["v"] for p in equity_curve] + for i in range(1, len(eq_vals)): + if eq_vals[i - 1] > 0: + returns.append(math.log(eq_vals[i] / eq_vals[i - 1])) + + wf_consistency = 0.5 if stats.sharpe > 0 else 0.0 + + verdict = QuantVerdict( + observed_sharpe=stats.sharpe, + wf_consistency=wf_consistency, + n_trials=4, + n_periods=max(total_trades, 1), + positive_regimes=1 if stats.sharpe > 0 else 0, + ).evaluate() + + vpin_curve = [] + if isinstance(self._maker, VPINGatedASMaker): + vpin_curve = self._maker._vpins[-200:] + + result = { + "strategy": self._maker_type, + "coin": coin.upper(), + "start_date": start_date, + "end_date": end_date, + "duration_hours": round(duration_hours, 2), + "n_events": len(events), + "total_trades": total_trades, + "bid_fills": stats.bid_fills, + "ask_fills": stats.ask_fills, + "cancels": stats.cancels, + "toxic_fills": stats.toxic_fills, + "adverse_rate": stats.adverse_rate, + "avg_spread_bps": stats.avg_spread_bps, + "max_inventory": stats.max_inventory, + "max_drawdown_pct": stats.max_drawdown, + "sharpe": stats.sharpe, + "sortino": stats.sortino, + "uptime_pct": stats.uptime_pct, + "pnl_breakdown": { + "spread_capture": breakdown.spread_capture, + "inventory_pnl": breakdown.inventory_pnl, + "maker_fees": breakdown.maker_fees, + "taker_fees": breakdown.taker_fees, + "funding_pnl": breakdown.funding_pnl, + "adverse_selection_cost": breakdown.adverse_selection_cost, + "gross_pnl": breakdown.gross_pnl, + "net_pnl": breakdown.net_pnl, + }, + "equity_curve": equity_curve, + "vpin_curve": vpin_curve, + "verdict": verdict["verdict"], + "dsr": verdict["deflated_sharpe"], + "psr": verdict["psr"], + "haircut_sharpe": verdict["haircut_sharpe"], + "maker_fee_pct": self._maker_fee_pct, + "taker_fee_pct": self._taker_fee_pct, + "maker_params": { + "gamma": self._maker.config.gamma, + "base_size": self._maker.config.base_size, + "max_inventory": self._maker.config.max_inventory, + "skew_factor": self._maker.config.skew_factor, + "vpin_threshold": self._vpin_threshold if self._maker_type == "vpin_as_mm" else None, + "vpin_alarm": self._vpin_alarm if self._maker_type == "vpin_as_mm" else None, + }, + "generated_at": datetime.now(timezone.utc).isoformat(), + } + + self._save_result(result) + return result + + def _empty_result(self, coin: str, start_date: str, end_date: str) -> dict: + return { + "strategy": self._maker_type, + "coin": coin.upper(), + "start_date": start_date, + "end_date": end_date, + "duration_hours": 0, + "n_events": 0, + "total_trades": 0, + "bid_fills": 0, + "ask_fills": 0, + "cancels": 0, + "toxic_fills": 0, + "adverse_rate": 0, + "max_drawdown_pct": 0, + "sharpe": 0, + "sortino": 0, + "pnl_breakdown": {}, + "equity_curve": [], + "vpin_curve": [], + "verdict": "INSUFFICIENT_DATA", + "error": "No events available", + "generated_at": datetime.now(timezone.utc).isoformat(), + } + + def _save_result(self, result: dict): + coin = result["coin"] + strategy = result["strategy"] + start = result["start_date"] + end = result["end_date"] + fname = f"{strategy}_{coin}_{start}_{end}_{datetime.now(timezone.utc).strftime('%Y%m%d-%H%M%S')}.json" + fpath = RESULTS_DIR / fname + with open(fpath, "w") as f: + json.dump(result, f, default=str) + logger.info("Saved result to %s", fpath) + + def print_summary(self, result: dict): + print("\n" + "=" * 60) + print(f" Tick Backtest — {result['strategy']} on {result['coin']}") + print(f" Period: {result['start_date']} → {result['end_date']}") + print(f" Duration: {result['duration_hours']}h | Events: {result['n_events']:,}") + print(f" Verdict: {result['verdict']}") + print("=" * 60) + print(f" Trades: {result['total_trades']} " + f"(bids: {result['bid_fills']}, asks: {result['ask_fills']})") + print(f" Cancels: {result['cancels']}") + print(f" Toxic fills: {result['toxic_fills']} ({result['adverse_rate']:.1%})") + print(f" Avg spread: {result['avg_spread_bps']:.1f} bps") + print(f" Max inventory: {result['max_inventory']:.6f}") + print(f" Max drawdown: {result['max_drawdown_pct']:.2f}%") + print(f" Sharpe: {result['sharpe']:.3f} Sortino: {result['sortino']:.3f}") + print(f" DSR: {result['dsr']:.4f} PSR: {result['psr']:.4f} " + f"Haircut: {result['haircut_sharpe']:.4f}") + print(f"\n PnL Breakdown:") + bd = result["pnl_breakdown"] + print(f" Spread capture: ${bd.get('spread_capture', 0):>10.4f}") + print(f" Inventory PnL: ${bd.get('inventory_pnl', 0):>10.4f}") + print(f" Maker fees: ${bd.get('maker_fees', 0):>10.4f}") + print(f" Taker fees: ${bd.get('taker_fees', 0):>10.4f}") + print(f" Funding PnL: ${bd.get('funding_pnl', 0):>10.4f}") + print(f" Adverse selection: ${bd.get('adverse_selection_cost', 0):>10.4f}") + print(f" " + "─" * 35) + print(f" Gross PnL: ${bd.get('gross_pnl', 0):>10.4f}") + print(f" Net PnL: ${bd.get('net_pnl', 0):>10.4f}") + print(f"\n Fee model: maker={result['maker_fee_pct']*100:.2f}% " + f"taker={result['taker_fee_pct']*100:.2f}%") + + +def cmd_tick_backtest(args): + runner = TickBacktestRunner( + data_dir=args.data_dir, + maker_type=args.maker, + gamma=args.gamma, + base_size=args.base_size, + max_inventory=args.max_inventory, + skew_factor=args.skew_factor, + vpin_threshold=args.vpin_threshold, + vpin_alarm=args.vpin_alarm, + maker_fee_pct=args.maker_fee / 100, + taker_fee_pct=args.taker_fee / 100, + adverse_selection_prob=args.adverse_prob, + cancel_after_ms=args.cancel_after_ms, + quote_refresh_ms=args.quote_refresh_ms, + seed=args.seed, + ) + + result = runner.run( + coin=args.coin, + start_date=args.start_date, + end_date=args.end_date, + ) + + runner.print_summary(result) + return result + + +if __name__ == "__main__": + logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s", datefmt="%H:%M:%S") + + p = argparse.ArgumentParser(description="Tick-level backtest runner") + p.add_argument("--coin", default="BTC") + p.add_argument("--data-dir", default="data/raw") + p.add_argument("--start-date", default="2026-08-01") + p.add_argument("--end-date", default="2026-08-07") + p.add_argument("--maker", default="as_mm", choices=["as_mm", "vpin_as_mm"]) + p.add_argument("--gamma", type=float, default=0.1) + p.add_argument("--base-size", type=float, default=0.001) + p.add_argument("--max-inventory", type=float, default=0.005) + p.add_argument("--skew-factor", type=float, default=0.5) + p.add_argument("--vpin-threshold", type=float, default=0.30) + p.add_argument("--vpin-alarm", type=float, default=0.50) + p.add_argument("--maker-fee", type=float, default=0.02, help="Maker fee in % (0.02 = 2bps)") + p.add_argument("--taker-fee", type=float, default=0.05, help="Taker fee in % (0.05 = 5bps)") + p.add_argument("--adverse-prob", type=float, default=0.15) + p.add_argument("--cancel-after-ms", type=float, default=5000.0) + p.add_argument("--quote-refresh-ms", type=float, default=2000.0) + p.add_argument("--seed", type=int, default=42) + args = p.parse_args() + + cmd_tick_backtest(args) diff --git a/cli.py b/cli.py index 75bd7cc..89fc597 100644 --- a/cli.py +++ b/cli.py @@ -105,6 +105,71 @@ def cmd_analyze(args): regime = funding_regime(rates, window_hours=24, n_samples_per_hour=1) print(f"Funding regime: {json.dumps(regime, indent=2, default=str)}") + elif args.channel == "markouts": + from microstructure.trades import compute_markouts, markout_summary + + l2_messages = read_range( + args.data_dir, + channel="l2book", + coin=args.coin.upper(), + start_date=args.start_date, + end_date=args.end_date, + ) + trade_messages = read_range( + args.data_dir, + channel="trades", + coin=args.coin.upper(), + start_date=args.start_date, + end_date=args.end_date, + ) + + trades = [] + mids = [] + times = [] + book_state = {} + for msg in sorted(l2_messages + trade_messages, key=lambda m: m.get("exchange_ts", 0) or 0): + ch = msg.get("channel", "") + payload = msg.get("payload", {}) + ts = msg.get("exchange_ts", 0) or 0 + if ch == "l2book": + levels = payload.get("levels", []) + if isinstance(levels, list) and len(levels) >= 2: + bids = [float(l["px"]) for l in levels[0] if float(l.get("sz", 0)) > 0] + asks = [float(l["px"]) for l in levels[1] if float(l.get("sz", 0)) > 0] + if bids and asks: + book_state["mid"] = (bids[0] + asks[0]) / 2 + elif ch == "trades": + px = float(payload.get("px", 0)) + if px > 0: + trades.append(payload) + mid = book_state.get("mid", px) + mids.append(mid) + times.append(ts) + + if not trades: + print("No trade data with L2 context available for markout analysis.") + return + + print(f"Analyzing {len(trades)} trades with L2 context...") + markouts = compute_markouts(trades, mids, times) + summary = markout_summary(markouts) + + print(f"\n{'─' * 70}") + print(f"{'Horizon':>10s} {'Buy Mean':>10s} {'Buy T-Stat':>10s} {'Buy N':>7s} " + f"{'Sell Mean':>10s} {'Sell T-Stat':>10s} {'Sell N':>7s}") + print(f"{'─' * 70}") + horizons = [100, 500, 1000, 5000, 10000, 30000, 60000] + for h in horizons: + b = summary.get("buy", {}).get(h, {}) + s = summary.get("sell", {}).get(h, {}) + print(f"{f'{h}ms':>10s} " + f"{b.get('mean_bps', 0):>10.2f} {b.get('t_stat', 0):>10.3f} {b.get('count', 0):>7d} " + f"{s.get('mean_bps', 0):>10.2f} {s.get('t_stat', 0):>10.3f} {s.get('count', 0):>7d}") + + print(f"\nBuy markout: + = price rises after buy (good for seller, bad for buyer)") + print(f"Sell markout: + = price falls after sell (good for buyer, bad for seller)") + print(f"t-stat > 2.0 = statistically significant predictive power") + else: print(f"Channel '{args.channel}' — raw dump:") for msg in messages[:5]: @@ -243,6 +308,197 @@ def cmd_backtest(args): print(f"\n{len(result['trades'])} trades") +def cmd_discover(args): + """Signal discovery — test microstructure signals against forward returns. + + For each L2 snapshot, computes WQI, OBI, VPIN, depth imbalance, and microprice. + Then measures how well each signal predicts mid-price movement at multiple horizons. + """ + from data.store import read_range + from microstructure.book import ( + mid_price, order_book_imbalance, depth_imbalance, microprice, spread_stats, depth_resiliency, + ) + from microstructure.trades import classify_lee_ready + from microstructure.toxicity import compute_vpin + from strategies.queue_imbalance import QueueImbalance + import numpy as np + + horizons = [int(h) for h in args.horizons.split(",")] + + l2_msgs = read_range( + args.data_dir, + channel="l2book", + coin=args.coin.upper(), + start_date=args.start_date, + end_date=args.end_date, + ) + trade_msgs = read_range( + args.data_dir, + channel="trades", + coin=args.coin.upper(), + start_date=args.start_date, + end_date=args.end_date, + ) + + if not l2_msgs or not trade_msgs: + print("No data available for signal discovery.") + return + + print(f"Loading {len(l2_msgs)} L2 messages and {len(trade_msgs)} trades for {args.coin}...") + + book_state = {"bids": {}, "asks": {}, "mid": 0.0, "ts": 0.0} + buy_vol = [] + sell_vol = [] + prev_bids = {} + prev_asks = {} + qi = QueueImbalance(depth_levels=10) + prev_mid = 0.0 + + events = [] + for msg in sorted(l2_msgs + trade_msgs, key=lambda m: m.get("exchange_ts", 0) or 0): + ch = msg.get("channel", "") + payload = msg.get("payload", {}) + ts = float(msg.get("exchange_ts", 0) or 0) / 1000.0 + events.append((ts, ch, payload)) + + events.sort(key=lambda e: e[0]) + + signals = [] + mids_series = [] + times_series = [] + + for ts, ch, payload in events: + if ch == "l2book": + bids = {} + asks = {} + levels = payload.get("levels", []) + if isinstance(levels, list) and len(levels) >= 2: + bid_list = [(float(l["px"]), float(l["sz"])) for l in levels[0] if float(l.get("sz", 0)) > 0] + ask_list = [(float(l["px"]), float(l["sz"])) for l in levels[1] if float(l.get("sz", 0)) > 0] + bids = dict(bid_list) + asks = dict(ask_list) + + if bids and asks: + book_state["bids"] = bids + book_state["asks"] = asks + mid = mid_price(bids, asks) + book_state["mid"] = mid + book_state["ts"] = ts + + obi = order_book_imbalance(bids, asks) + di = depth_imbalance(bids, asks) + mp = microprice(bids, asks) + ss = spread_stats(bids, asks) + dr = depth_resiliency(bids, asks) + + wqi = qi.compute_wqi(bid_list, ask_list) + + vpin_val = 0.0 + if buy_vol and sell_vol: + v = compute_vpin(buy_vol, sell_vol, n_buckets=50) + vpin_val = v.get("vpin_value", 0.0) + + bid_depth = sum(sz for _, sz in bid_list[:10]) + ask_depth = sum(sz for _, sz in ask_list[:10]) + total_depth = bid_depth + ask_depth + + signals.append({ + "ts": ts, + "mid": mid, + "obi": obi, + "wqi": wqi, + "vpin": vpin_val, + "depth_imbalance": di, + "microprice_ratio": mp / mid if mid > 0 else 1.0, + "spread_bps": ss["spread_bps"], + "bid_depth": bid_depth, + "ask_depth": ask_depth, + "depth_total": total_depth, + "resiliency": dr.get("resiliency", 0), + }) + mids_series.append(mid) + times_series.append(ts) + prev_bids = bid_list + prev_asks = ask_list + + elif ch == "trades": + px = float(payload.get("px", 0)) + sz = float(payload.get("sz", 0)) + if px > 0 and sz > 0: + side = classify_lee_ready(px, book_state["mid"]) + if side == "buy": + buy_vol.append(sz) + else: + sell_vol.append(sz) + + n = len(signals) + if n < 50: + print("Too few L2 snapshots for signal discovery. Need more data.") + return + + print(f"Signal discovery on {n} L2 snapshots across {horizons}ms horizons...") + + signal_names = ["obi", "wqi", "vpin", "depth_imbalance", "spread_bps", "resiliency"] + horizons = sorted(horizons) + + print(f"\n{'─' * 80}") + print(f"{'Signal':>18s} ", end="") + for h in horizons: + print(f" {'t-{h}ms':>10s}", end="") + print(f" {'r²':>8s}") + print(f"{'─' * 80}") + + for sig_name in signal_names: + sig_vals = [s.get(sig_name, 0) for s in signals] + print(f"{sig_name:>18s} ", end="") + + for horizon in horizons: + t_stats = [] + for i in range(n - 1): + future_idx = i + for j in range(i + 1, min(n, len(times_series))): + if times_series[j] - times_series[i] >= horizon / 1000.0: + future_idx = j + break + if future_idx > i and mids_series[i] > 0: + forward_return = (mids_series[future_idx] - mids_series[i]) / mids_series[i] * 10000 + if abs(forward_return) < 500: + zipped = list(zip(sig_vals, [forward_return] * len(sig_vals))) + t_tuple = zipped[i] if i < len(zipped) else None + + t_stats = [] + for i in range(n - 1): + future_idx = i + for j in range(i + 1, n): + if times_series[j] - times_series[i] >= horizon / 1000.0: + future_idx = j + break + if future_idx > i and mids_series[i] > 0: + forward_return = (mids_series[future_idx] - mids_series[i]) / mids_series[i] * 10000 + signal_val = sig_vals[i] + if abs(forward_return) < 500 and abs(signal_val) < 100: + t_stats.append((signal_val, forward_return)) + + if len(t_stats) >= 10: + xs = np.array([t[0] for t in t_stats]) + ys = np.array([t[1] for t in t_stats]) + r = np.corrcoef(xs, ys)[0, 1] if len(xs) > 1 else 0 + t_stat = r * np.sqrt(len(t_stats) - 2) / np.sqrt(1 - r * r) if abs(r) < 1 else 0 + print(f" {t_stat:>10.3f}", end="") + else: + print(f" {'N/A':>10s}", end="") + + if len(t_stats) >= 10: + xs = np.array([t[0] for t in t_stats]) + ys = np.array([t[1] for t in t_stats]) + r_sq = np.corrcoef(xs, ys)[0, 1] ** 2 if len(xs) > 1 else 0 + print(f" {r_sq:>8.4f}") + else: + print() + + print(f"\nPipeline ready. Run 'python -m cli tick' to backtest strategies on this data.") + + def main(): import argparse p = argparse.ArgumentParser(description="FTDT Quant Lab CLI") @@ -259,7 +515,7 @@ def main(): # analyze pa = sp.add_parser("analyze", help="Run microstructure analytics") pa.add_argument("--data-dir", default="data/raw") - pa.add_argument("--channel", default="l2book", choices=["l2book", "trades", "funding", "mark", "open_interest", "liquidation"]) + pa.add_argument("--channel", default="l2book", choices=["l2book", "trades", "funding", "mark", "open_interest", "liquidation", "markouts"]) pa.add_argument("--coin", default="BTC") pa.add_argument("--start-date", default="2026-08-01") pa.add_argument("--end-date", default="2026-08-07") @@ -294,6 +550,34 @@ def main(): pb.add_argument("--interval", default="1h") pb.add_argument("--limit", type=int, default=500) + # tick + pt = sp.add_parser("tick", help="Tick-level backtest (Parquet L2+trade replay)") + pt.add_argument("--coin", default="BTC") + pt.add_argument("--data-dir", default="data/raw") + pt.add_argument("--start-date", default="2026-08-01") + pt.add_argument("--end-date", default="2026-08-07") + pt.add_argument("--maker", default="as_mm", choices=["as_mm", "vpin_as_mm"]) + pt.add_argument("--gamma", type=float, default=0.1) + pt.add_argument("--base-size", type=float, default=0.001) + pt.add_argument("--max-inventory", type=float, default=0.005) + pt.add_argument("--skew-factor", type=float, default=0.5) + pt.add_argument("--vpin-threshold", type=float, default=0.30) + pt.add_argument("--vpin-alarm", type=float, default=0.50) + pt.add_argument("--maker-fee", type=float, default=0.02, help="Maker fee in %") + pt.add_argument("--taker-fee", type=float, default=0.05, help="Taker fee in %") + pt.add_argument("--adverse-prob", type=float, default=0.15) + pt.add_argument("--cancel-after-ms", type=float, default=5000.0) + pt.add_argument("--quote-refresh-ms", type=float, default=2000.0) + pt.add_argument("--seed", type=int, default=42) + + # discover + pd = sp.add_parser("discover", help="Signal discovery — test microstructure signals against forward returns") + pd.add_argument("--data-dir", default="data/raw") + pd.add_argument("--coin", default="BTC") + pd.add_argument("--start-date", default="2026-08-01") + pd.add_argument("--end-date", default="2026-08-07") + pd.add_argument("--horizons", default="100,500,1000,5000,10000", help="Comma-separated ms horizons") + args = p.parse_args() import json as _json @@ -309,6 +593,11 @@ def main(): cmd_run(args) elif args.command == "backtest": cmd_backtest(args) + elif args.command == "tick": + from backtests.tick_runner import cmd_tick_backtest + cmd_tick_backtest(args) + elif args.command == "discover": + cmd_discover(args) if __name__ == "__main__": diff --git a/sim/engine.py b/sim/engine.py index 0b52fe9..f8092d5 100644 --- a/sim/engine.py +++ b/sim/engine.py @@ -60,10 +60,31 @@ class SimConfig: scenario: ScenarioConfig = field(default_factory=ScenarioConfig) # Simulation behavior - cancel_after_ms: float = 5000.0 # cancel and re-quote every N ms - quote_refresh_ms: float = 2000.0 # refresh quotes every N ms + cancel_after_ms: float = 5000.0 + quote_refresh_ms: float = 2000.0 seed: int | None = None + @classmethod + def from_fee_tier( + cls, + vip_tier: int = 0, + staking_tier: str = "none", + maker_rebate_tier: int = 0, + **kwargs, + ) -> "SimConfig": + from config.fee_tiers import get_perp_fees, PERPS_TIERS, STAKING_TIERS + + maker_fee = get_perp_fees(vip_tier, staking_tier, "maker", maker_rebate_tier) + taker_fee = get_perp_fees(vip_tier, staking_tier, "taker", maker_rebate_tier) + tier_name = PERPS_TIERS[vip_tier]["name"] + staking_name = STAKING_TIERS.get(staking_tier, STAKING_TIERS["none"])["name"] + + return cls( + maker_fee_pct=maker_fee, + taker_fee_pct=taker_fee, + **kwargs, + ) + class SimulationEngine: """Event-driven market-making simulator. diff --git a/sim/fills.py b/sim/fills.py index e189a89..222c633 100644 --- a/sim/fills.py +++ b/sim/fills.py @@ -180,3 +180,118 @@ def adverse_selection_intensity( "total_cost_bps": round(sum(adverse_costs), 2), "n_fills": len(fills), } + + +class QueueAwareFillModel: + """Realistic queue-priority fill model for paper trading. + + Unlike random fills, this models whether an aggressor trade at a given + price level would exhaust the queue ahead of our order. An order fills + only when the total aggressor volume at that price exceeds the volume + of orders ahead of ours in the FIFO queue. + + Usage: + model = QueueAwareFillModel() + fill_info = model.check_fill( + aggressor_side="buy", + agg_size=0.005, + our_price=50000.0, + our_size=0.001, + depth_ahead=0.002, + ) + """ + + def __init__(self, base_fill_prob: float = 0.20): + self._base_fill_prob = base_fill_prob + self._fill_count: int = 0 + self._skip_count: int = 0 + + def check_fill( + self, + aggressor_side: str, + agg_size: float, + agg_price: float, + our_price: float, + our_size: float, + depth_ahead: float, + ) -> dict: + """Determine if a market order would fill our limit order. + + Args: + aggressor_side: "buy" (market buy hits asks) or "sell" (market sell hits bids) + agg_size: size of the aggressor trade + agg_price: price of the aggressor trade + our_price: our limit order price + our_size: our order size + depth_ahead: total size of orders ahead of ours in the queue at this price + + Returns: + dict with filled (bool), fill_size, reason + """ + price_match = False + if aggressor_side == "buy" and agg_price >= our_price: + price_match = True + elif aggressor_side == "sell" and agg_price <= our_price: + price_match = True + + if not price_match: + return {"filled": False, "fill_size": 0.0, "reason": "price_not_crossed"} + + remaining_after_queue = agg_size - depth_ahead + if remaining_after_queue <= 0: + self._skip_count += 1 + return {"filled": False, "fill_size": 0.0, "reason": "queue_not_reached"} + + fill_size = min(our_size, remaining_after_queue) + self._fill_count += 1 + + return { + "filled": True, + "fill_size": round(fill_size, 8), + "fill_ratio": round(fill_size / our_size, 4), + "reason": f"reached_queue_pos", + "depth_consumed": round(depth_ahead + fill_size, 8), + } + + def estimate_depth_ahead( + self, + our_price: float, + our_side: str, + best_bid: float, + best_ask: float, + bid_depth: float, + ask_depth: float, + ) -> float: + """Estimate the volume ahead of our order at a price level. + + This is a heuristic since HL doesn't expose queue position. We estimate + based on whether we're at the best level and how much total depth is there. + """ + at_best = ( + (our_side == "bid" and our_price >= best_bid) or + (our_side == "ask" and our_price <= best_ask) + ) + + if not at_best: + return float("inf") + + if our_side == "bid": + return bid_depth * 0.5 + else: + return ask_depth * 0.5 + + @property + def fill_count(self) -> int: + return self._fill_count + + @property + def skip_count(self) -> int: + return self._skip_count + + def fill_rate(self) -> float: + total = self._fill_count + self._skip_count + return self._fill_count / total if total > 0 else 0.0 + + def reset(self): + self._fill_count = 0 + self._skip_count = 0 diff --git a/strategies/wqi_predictor.py b/strategies/wqi_predictor.py new file mode 100644 index 0000000..50519ab --- /dev/null +++ b/strategies/wqi_predictor.py @@ -0,0 +1,342 @@ +""" +WQI Predictor — Queue Imbalance Directional Strategy. + +Uses the Weighted Queue Imbalance (WQI) from `strategies/queue_imbalance.py` +to predict short-term price direction. The strategy enters when WQI z-score +is extreme and exits on timeout, reversal, or stop-loss. + +Based on the Stoikov-Sağlam and Cont frameworks for order book dynamics. + +Entry logic: + - WQI z-score > z_entry (default 2.0) AND WQI > wqi_threshold (0.15) → BUY + - WQI z-score < -z_entry AND WQI < -wqi_threshold → SELL + - Gate: adverse selection score < max_adverse (0.3) + +Exit logic: + - WQI crosses zero (microstructure reversal) + - Hold time exceeds max_hold_seconds + - Stop-loss: price moves > stop_loss_bps against position + - Profit target: price moves > take_profit_bps in favor + +Usage (tick-level backtest): + predictor = WQIPredictor(z_entry=2.0, max_hold_seconds=30) + for l2_snapshot, trade in replay: + signal = predictor.update(bids, asks, mid_price, prev_bids, prev_asks, prev_mid) + if signal["action"] in ("BUY", "SELL"): + execute_trade(signal) + +Usage (live paper trading): + predictor = WQIPredictor(adverse_trades=True) + signal = predictor.feed_signal(bids, asks, mid_price, trade_side, trade_size) +""" + +from __future__ import annotations + +import math +import time +from collections import deque +from typing import Optional + +from strategies.queue_imbalance import QueueImbalance + + +class WQIPredictor: + """Queue-imbalance-based directional trading strategy. + + The predictive signal: when depth at the top of the book systematically + shifts to one side, the mid-price tends to move in that direction. + This strategy captures that by entering when the deviation is extreme + relative to recent history and exiting when it normalizes. + """ + + def __init__( + self, + z_entry: float = 2.0, + z_exit: float = 0.5, + wqi_threshold: float = 0.15, + max_adverse: float = 0.3, + max_hold_seconds: float = 30.0, + stop_loss_bps: float = 5.0, + take_profit_bps: float = 10.0, + size: float = 0.001, + depth_levels: int = 10, + fee_model: str = "taker", + ): + self._z_entry = z_entry + self._z_exit = z_exit + self._wqi_threshold = wqi_threshold + self._max_adverse = max_adverse + self._max_hold_seconds = max_hold_seconds + self._stop_loss_bps = stop_loss_bps + self._take_profit_bps = take_profit_bps + self._size = size + self._fee_model = fee_model + + self._qi = QueueImbalance(depth_levels=depth_levels) + self._mid_prices: deque[float] = deque(maxlen=200) + self._signals: deque[dict] = deque(maxlen=100) + + self._position: int = 0 + self._entry_price: float = 0.0 + self._entry_time: float = 0.0 + self._entry_wqi: float = 0.0 + + self._trades: list[dict] = [] + + def feed_signal( + self, + bids: list, + asks: list, + mid_price: float, + prev_bids: Optional[list] = None, + prev_asks: Optional[list] = None, + prev_mid: float = 0.0, + timestamp: Optional[float] = None, + ) -> dict: + if timestamp is None: + timestamp = time.time() + + analysis = self._qi.analyze(bids, asks, mid_price, prev_bids, prev_asks, prev_mid) + self._mid_prices.append(mid_price) + + wqi = analysis["wqi"] + z_score = analysis["z_score"] + adverse = analysis["adverse_selection"] + net_flow = analysis["net_flow"] + + action = "HOLD" + reason = "" + + if self._position == 0: + if z_score > self._z_entry and wqi > self._wqi_threshold and adverse < self._max_adverse: + action = "BUY" + reason = f"wqi_z={z_score:.2f}_wqi={wqi:.3f}_adv={adverse:.3f}" + elif z_score < -self._z_entry and wqi < -self._wqi_threshold and adverse < self._max_adverse: + action = "SELL" + reason = f"wqi_z={z_score:.2f}_wqi={wqi:.3f}_adv={adverse:.3f}" + else: + hold_time = timestamp - self._entry_time + pnl_bps = (mid_price / self._entry_price - 1) * 10000 + if self._position == -1: + pnl_bps = -pnl_bps + + if abs(z_score) < self._z_exit: + action = "EXIT" + reason = f"z_cross={z_score:.2f}" + elif hold_time >= self._max_hold_seconds: + action = "EXIT" + reason = f"timeout_{hold_time:.0f}s" + elif pnl_bps <= -self._stop_loss_bps: + action = "EXIT" + reason = f"stop_loss_{pnl_bps:.1f}bps" + elif pnl_bps >= self._take_profit_bps: + action = "EXIT" + reason = f"take_profit_{pnl_bps:.1f}bps" + + if action in ("BUY", "SELL"): + self._position = 1 if action == "BUY" else -1 + self._entry_price = mid_price + self._entry_time = timestamp + self._entry_wqi = wqi + elif action == "EXIT": + pnl_bps = (mid_price / self._entry_price - 1) * 10000 + if self._position == 1: + gross_pnl = self._size * (mid_price - self._entry_price) + else: + gross_pnl = self._size * (self._entry_price - mid_price) + fee = self._size * mid_price * (0.00045 if self._fee_model == "taker" else 0.00015) + net_pnl = gross_pnl - fee + + self._trades.append({ + "entry_price": round(self._entry_price, 2), + "exit_price": round(mid_price, 2), + "side": "BUY" if self._position == 1 else "SELL", + "size": self._size, + "pnl_bps": round(pnl_bps, 2), + "gross_pnl": round(gross_pnl, 4), + "fee": round(fee, 4), + "net_pnl": round(net_pnl, 4), + "hold_seconds": round(timestamp - self._entry_time, 2), + "entry_wqi": round(self._entry_wqi, 4), + "exit_wqi": round(wqi, 4), + "reason": reason, + }) + self._position = 0 + self._entry_price = 0.0 + self._entry_time = 0.0 + self._entry_wqi = 0.0 + + signal = { + "action": action, + "wqi": wqi, + "z_score": z_score, + "adverse": adverse, + "net_flow": net_flow, + "mid_price": mid_price, + "position": self._position, + "reason": reason, + "timestamp": timestamp, + } + self._signals.append(signal) + return signal + + @property + def position(self) -> int: + return self._position + + @property + def trades(self) -> list[dict]: + return self._trades + + @property + def recent_signals(self) -> list[dict]: + return list(self._signals) + + def summary(self) -> dict: + if not self._trades: + return { + "total_trades": 0, + "win_rate": 0.0, + "total_net_pnl": 0.0, + "avg_net_pnl": 0.0, + "avg_hold_seconds": 0.0, + "max_profit_bps": 0.0, + "max_loss_bps": 0.0, + } + + wins = sum(1 for t in self._trades if t["net_pnl"] > 0) + total_net = sum(t["net_pnl"] for t in self._trades) + gross_pnls = [t["pnl_bps"] for t in self._trades] + hold_times = [t["hold_seconds"] for t in self._trades] + + return { + "total_trades": len(self._trades), + "win_rate": round(wins / len(self._trades), 3), + "total_net_pnl": round(total_net, 4), + "avg_net_pnl": round(total_net / len(self._trades), 4), + "avg_pnl_bps": round(sum(gross_pnls) / len(gross_pnls), 2), + "avg_hold_seconds": round(sum(hold_times) / len(hold_times), 2), + "max_profit_bps": round(max(gross_pnls), 2), + "max_loss_bps": round(min(gross_pnls), 2), + } + + def reset(self): + self._position = 0 + self._entry_price = 0.0 + self._entry_time = 0.0 + self._entry_wqi = 0.0 + self._trades.clear() + self._signals.clear() + self._mid_prices.clear() + self._qi = QueueImbalance(depth_levels=self._qi.depth_levels) + + +def run_wqi_backtest( + tick_data_path: str, + coin: str = "BTC", + z_entry: float = 2.0, + size: float = 0.001, + max_hold: float = 30.0, + stop_loss: float = 5.0, + take_profit: float = 10.0, + taker_fee_pct: float = 0.05, +) -> dict: + """Run WQI predictor backtest on stored tick data. + + Args: + tick_data_path: path to parquet data directory + coin: instrument + z_entry: WQI z-score entry threshold + size: trade size in BTC + max_hold: max hold time in seconds + stop_loss: stop loss in bps + take_profit: take profit in bps + taker_fee_pct: taker fee in % (0.05 = 5bps) + + Returns dict with trades, PnL, and metrics. + """ + import json + from data.store import read_range + from microstructure.book import mid_price + + predictor = WQIPredictor( + z_entry=z_entry, + max_hold_seconds=max_hold, + stop_loss_bps=stop_loss, + take_profit_bps=take_profit, + size=size, + fee_model="taker", + ) + + trade_msgs = read_range(tick_data_path, channel="l2book", coin=coin, + start_date="2026-01-01", end_date="2030-01-01") + + prev_bids = None + prev_asks = None + prev_mid = 0.0 + + for msg in trade_msgs: + payload = msg.get("payload", {}) + levels = payload.get("levels", []) + if not isinstance(levels, list) or len(levels) < 2: + continue + + bids = [(float(l["px"]), float(l["sz"])) for l in levels[0] if float(l.get("sz", 0)) > 0] + asks = [(float(l["px"]), float(l["sz"])) for l in levels[1] if float(l.get("sz", 0)) > 0] + if not bids or not asks: + continue + + current_mid = mid_price(dict(bids), dict(asks)) + if current_mid <= 0: + continue + + ts = msg.get("local_ts", time.time()) + + signal = predictor.feed_signal( + bids, asks, current_mid, + prev_bids, prev_asks, prev_mid, + timestamp=ts, + ) + + prev_bids = bids + prev_asks = asks + prev_mid = current_mid + + return predictor.summary() + + +def run_wqi_live_signal(bids, asks, mid_market, state: dict) -> dict: + """Stateless WQI signal for live trading integration. + + Args: + bids: list of [price, size] + asks: list of [price, size] + mid_market: current mid price + state: dict with 'prev_bids', 'prev_asks', 'prev_mid', 'position', + 'entry_price', 'entry_time' + + Returns signal dict with action and state updates. + """ + predictor = WQIPredictor() + predictor._position = state.get("position", 0) + predictor._entry_price = state.get("entry_price", 0.0) + predictor._entry_time = state.get("entry_time", time.time()) + + signal = predictor.feed_signal( + bids, asks, mid_market, + prev_bids=state.get("prev_bids"), + prev_asks=state.get("prev_asks"), + prev_mid=state.get("prev_mid", 0.0), + ) + + return { + **signal, + "state_update": { + "prev_bids": bids, + "prev_asks": asks, + "prev_mid": mid_market, + "position": predictor.position, + "entry_price": predictor._entry_price, + "entry_time": predictor._entry_time, + }, + } diff --git a/tests/test_tick_backtest.py b/tests/test_tick_backtest.py new file mode 100644 index 0000000..33f0728 --- /dev/null +++ b/tests/test_tick_backtest.py @@ -0,0 +1,228 @@ +""" +Tests for tick-level backtest runner and queue-aware fill model. +""" +from sim.fills import QueueAwareFillModel, FillModelConfig, FillSimulator +from sim.maker import MakerConfig +from sim.engine import SimConfig + + +class TestQueueAwareFillModel: + def test_price_not_crossed(self): + qm = QueueAwareFillModel() + result = qm.check_fill( + aggressor_side="buy", + agg_size=0.01, + agg_price=49900.0, + our_price=50000.0, + our_size=0.001, + depth_ahead=0.0, + ) + assert not result["filled"] + assert result["reason"] == "price_not_crossed" + + def test_fills_when_price_crossed_and_no_queue_ahead(self): + qm = QueueAwareFillModel() + result = qm.check_fill( + aggressor_side="buy", + agg_size=0.01, + agg_price=50005.0, + our_price=50000.0, + our_size=0.001, + depth_ahead=0.0, + ) + assert result["filled"] + assert result["fill_size"] == 0.001 + + def test_does_not_fill_when_queue_not_reached(self): + qm = QueueAwareFillModel() + result = qm.check_fill( + aggressor_side="buy", + agg_size=0.001, + agg_price=50005.0, + our_price=50000.0, + our_size=0.001, + depth_ahead=0.005, + ) + assert not result["filled"] + assert result["reason"] == "queue_not_reached" + + def test_partial_fill(self): + qm = QueueAwareFillModel() + result = qm.check_fill( + aggressor_side="sell", + agg_size=0.005, + agg_price=49990.0, + our_price=50000.0, + our_size=0.003, + depth_ahead=0.002, + ) + assert result["filled"] + assert result["fill_size"] == 0.003 + + def test_sell_fill_price_match(self): + qm = QueueAwareFillModel() + result = qm.check_fill( + aggressor_side="sell", + agg_size=0.01, + agg_price=49990.0, + our_price=50000.0, + our_size=0.001, + depth_ahead=0.0, + ) + assert result["filled"] + + def test_estimate_depth_ahead_at_best(self): + qm = QueueAwareFillModel() + depth = qm.estimate_depth_ahead( + our_price=50000.0, + our_side="bid", + best_bid=50000.0, + best_ask=50002.0, + bid_depth=2.0, + ask_depth=1.0, + ) + assert depth == 1.0 + + def test_estimate_depth_ahead_not_at_best(self): + qm = QueueAwareFillModel() + depth = qm.estimate_depth_ahead( + our_price=49999.0, + our_side="bid", + best_bid=50000.0, + best_ask=50002.0, + bid_depth=2.0, + ask_depth=1.0, + ) + assert depth == float("inf") + + def test_fill_rate_tracking(self): + qm = QueueAwareFillModel() + qm.check_fill("buy", 0.01, 50005.0, 50000.0, 0.001, 0.0) + qm.check_fill("buy", 0.001, 50005.0, 50000.0, 0.001, 0.005) + qm.check_fill("buy", 0.01, 50005.0, 50000.0, 0.001, 0.0) + assert qm.fill_count == 2 + assert qm.skip_count == 1 + assert qm.fill_rate() == 2 / 3 + + +class TestSimConfigFeeTier: + def test_from_fee_tier_default(self): + cfg = SimConfig.from_fee_tier(vip_tier=0) + assert cfg.maker_fee_pct == 0.00015 + assert cfg.taker_fee_pct == 0.00045 + + def test_from_fee_tier_vip2(self): + cfg = SimConfig.from_fee_tier(vip_tier=2) + assert cfg.maker_fee_pct == 0.00008 + assert cfg.taker_fee_pct == 0.00035 + + def test_from_fee_tier_with_staking(self): + cfg = SimConfig.from_fee_tier(vip_tier=0, staking_tier="gold") + assert cfg.maker_fee_pct < 0.00015 + assert cfg.taker_fee_pct < 0.00045 + + def test_from_fee_tier_custom_params(self): + cfg = SimConfig.from_fee_tier( + vip_tier=0, + max_inventory=0.01, + initial_equity=50000.0, + ) + assert cfg.max_inventory == 0.01 + assert cfg.initial_equity == 50000.0 + + +class TestVPINGatedASMaker: + def test_default_allows_quoting(self): + from backtests.tick_runner import VPINGatedASMaker + maker = VPINGatedASMaker(MakerConfig()) + assert maker.allowed_to_quote() == (True, 1.0) + + def test_alarm_blocks_quoting(self): + from backtests.tick_runner import VPINGatedASMaker + maker = VPINGatedASMaker(MakerConfig(), vpin_threshold=0.3, vpin_alarm=0.5) + maker._current_vpin = 0.55 + assert maker.allowed_to_quote() == (False, 0.0) + + def test_threshold_reduces_size(self): + from backtests.tick_runner import VPINGatedASMaker + maker = VPINGatedASMaker(MakerConfig(), vpin_threshold=0.3, vpin_alarm=0.5) + maker._current_vpin = 0.35 + allowed, size_mult = maker.allowed_to_quote() + assert allowed + assert 0 < size_mult < 1.0 + + def test_returns_none_when_not_allowed(self): + from backtests.tick_runner import VPINGatedASMaker + maker = VPINGatedASMaker(MakerConfig(), vpin_threshold=0.3, vpin_alarm=0.5) + maker._current_vpin = 0.55 + maker.observe(100000.0) + q = maker.quote(100000.0, 0.0, 0.0) + assert q is None + + def test_returns_quote_when_allowed(self): + from backtests.tick_runner import VPINGatedASMaker + maker = VPINGatedASMaker(MakerConfig()) + maker.observe(100000.0) + maker.observe(100100.0) + maker.observe(100050.0) + q = maker.quote(100000.0, 0.0, 0.0) + assert q is not None + assert q.bid < q.ask + assert q.bid_size > 0 + + +class TestWQIPredictor: + def test_initial_state(self): + from strategies.wqi_predictor import WQIPredictor + wqi = WQIPredictor() + assert wqi.position == 0 + assert len(wqi.trades) == 0 + + def test_no_signal_with_balanced_book(self): + from strategies.wqi_predictor import WQIPredictor + wqi = WQIPredictor() + bids = [(100.0, 1.0), (99.0, 1.0)] + asks = [(102.0, 1.0), (103.0, 1.0)] + signal = wqi.feed_signal(bids, asks, 101.0) + assert signal["action"] == "HOLD" + + def test_buy_signal_with_bid_heavy_book(self): + from strategies.wqi_predictor import WQIPredictor + wqi = WQIPredictor(z_entry=0.5, wqi_threshold=0.05) + for _ in range(25): + wqi.feed_signal([(100.0, 1.0), (99.0, 1.0)], [(102.0, 1.0), (103.0, 1.0)], 101.0) + bids = [(100.0, 10.0), (99.0, 5.0)] + asks = [(102.0, 1.0)] + signal = wqi.feed_signal(bids, asks, 101.0) + if signal["action"] == "BUY": + assert wqi.position != 0 + + def test_exit_on_timeout(self): + from strategies.wqi_predictor import WQIPredictor + import time + wqi = WQIPredictor(z_entry=0.01, z_exit=999.0, wqi_threshold=0.01, + max_hold_seconds=0.001) + for _ in range(25): + wqi.feed_signal([(100.0, 1.0), (99.0, 1.0)], [(102.0, 1.0), (103.0, 1.0)], 101.0) + bids = [(100.0, 10.0), (99.0, 5.0)] + asks = [(102.0, 1.0)] + wqi.feed_signal(bids, asks, 101.0) + time.sleep(0.01) + signal = wqi.feed_signal(bids, asks, 101.0) + assert signal["action"] in ("HOLD", "EXIT") + + def test_summary_returns_zero_for_no_trades(self): + from strategies.wqi_predictor import WQIPredictor + wqi = WQIPredictor() + s = wqi.summary() + assert s["total_trades"] == 0 + + def test_reset_clears_state(self): + from strategies.wqi_predictor import WQIPredictor + wqi = WQIPredictor() + bids = [(100.0, 1.0)] + asks = [(102.0, 1.0)] + wqi.feed_signal(bids, asks, 101.0) + wqi.reset() + assert wqi.position == 0 + assert len(wqi.trades) == 0