feat: Phase 2 — microstructure analytics + 81 tests
New microstructure/ module with pure-function analytics: microstructure/book.py: microprice() — depth-weighted mid price mid_price() — simple bid/ask midpoint order_book_imbalance() — ranged [-1, 1] volume skew depth_imbalance() — imbalance at fixed price distance spread_stats() — spread, spread_bps, mid, bid, ask depth_resiliency() — bid/ask volume within impact radius queue_depletion_prob() — Poisson fill probability at level batch_book_stats() — aggregate stats across snapshots microstructure/trades.py: classify_lee_ready() — Lee-Ready aggressor classification classify_bulk_lee_ready() — batch classification with mids/bids/asks compute_markouts() — forward mid-price change at configurable horizons markout_summary() — mean/std/t-stat per side per horizon trade_volume_profile() — size bucket distribution trade_arrival_rate() — rolling trades/sec with burst detection microstructure/toxicity.py: compute_vpin() — volume-synchronized informed trading probability compute_vpin_time_series() — rolling VPIN with alarm threshold fill_toxicity() — adverse price movement post-trade adverse_selection_ratio() — per-side adverse selection liquidation_clustering() — cluster detection in liquidation events microstructure/funding.py: funding_regime() — classify regime (neutral/positive/negative/high) funding_predictability() — AR(1) autocorrelation analysis funding_carry_pnl() — cumulative carry PnL estimation basis_spread() — perp premium over spot (bps) basis_convergence_speed() — mean-reversion half-life via AR(1) microstructure/signals.py: composite_signal() — weighted OBI + trade + VPIN + funding signal SignalPipeline — stateful pipeline accumulating book/trade updates detect_hft_regime() — regime classifier for HFT strategy selection Bug fixes in Phase 1: - data/latency.py: proper linear-interpolation percentiles - data/normalizer.py: UTC timezone for naive datetimes - data/normalizer.py: detect_sequence_gap returns gap-1 (missing count) - microstructure/toxicity.py: consistent vpin_value key in compute_vpin 81 tests across 4 test files (store, normalizer, latency, microstructure)
This commit is contained in:
+12
-4
@@ -80,11 +80,19 @@ class LatencyTracker:
|
||||
return {"p50": 0, "p90": 0, "p95": 0, "p99": 0, "count": 0}
|
||||
vals = sorted(v for _, v in buffer)
|
||||
n = len(vals)
|
||||
|
||||
def _pct(p: float) -> float:
|
||||
k = (n - 1) * p / 100
|
||||
lo = int(k)
|
||||
hi = min(lo + 1, n - 1)
|
||||
frac = k - lo
|
||||
return vals[lo] + frac * (vals[hi] - vals[lo])
|
||||
|
||||
return {
|
||||
"p50": round(vals[int(n * 0.50)], 2),
|
||||
"p90": round(vals[int(n * 0.90)], 2),
|
||||
"p95": round(vals[int(n * 0.95)], 2),
|
||||
"p99": round(vals[int(n * 0.99)], 2),
|
||||
"p50": round(_pct(50), 2),
|
||||
"p90": round(_pct(90), 2),
|
||||
"p95": round(_pct(95), 2),
|
||||
"p99": round(_pct(99), 2),
|
||||
"max": round(vals[-1], 2),
|
||||
"count": n,
|
||||
}
|
||||
|
||||
+3
-1
@@ -39,6 +39,8 @@ def normalize_timestamp(ts, source: str = "hl") -> int:
|
||||
return _parse_iso_ms(ts)
|
||||
|
||||
if isinstance(ts, datetime):
|
||||
if ts.tzinfo is None:
|
||||
ts = ts.replace(tzinfo=timezone.utc)
|
||||
return int(ts.timestamp() * 1000)
|
||||
|
||||
return int(time.time() * 1000)
|
||||
@@ -120,5 +122,5 @@ def detect_sequence_gap(
|
||||
if diff == 1:
|
||||
return 0
|
||||
if diff > 1:
|
||||
return min(diff, max_gap * 100) # cap reporting size
|
||||
return min(diff - 1, max_gap * 100)
|
||||
return -1 # duplicate or reset
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
"""
|
||||
Microstructure analytics — book, trade, toxicity, funding, and signal composition.
|
||||
|
||||
Pure functions that take market data arrays and return computed metrics.
|
||||
"""
|
||||
|
||||
from microstructure.book import (
|
||||
microprice,
|
||||
mid_price,
|
||||
order_book_imbalance,
|
||||
depth_imbalance,
|
||||
spread_stats,
|
||||
depth_resiliency,
|
||||
queue_depletion_prob,
|
||||
batch_book_stats,
|
||||
)
|
||||
|
||||
from microstructure.trades import (
|
||||
classify_lee_ready,
|
||||
classify_bulk_lee_ready,
|
||||
compute_markouts,
|
||||
markout_summary,
|
||||
trade_volume_profile,
|
||||
trade_arrival_rate,
|
||||
)
|
||||
|
||||
from microstructure.toxicity import (
|
||||
compute_vpin,
|
||||
compute_vpin_time_series,
|
||||
fill_toxicity,
|
||||
adverse_selection_ratio,
|
||||
liquidation_clustering,
|
||||
)
|
||||
|
||||
from microstructure.funding import (
|
||||
funding_regime,
|
||||
funding_predictability,
|
||||
funding_carry_pnl,
|
||||
basis_spread,
|
||||
basis_convergence_speed,
|
||||
)
|
||||
|
||||
from microstructure.signals import (
|
||||
composite_signal,
|
||||
SignalPipeline,
|
||||
detect_hft_regime,
|
||||
)
|
||||
|
||||
__all__ = [
|
||||
"microprice", "mid_price", "order_book_imbalance", "depth_imbalance",
|
||||
"spread_stats", "depth_resiliency", "queue_depletion_prob", "batch_book_stats",
|
||||
"classify_lee_ready", "classify_bulk_lee_ready", "compute_markouts",
|
||||
"markout_summary", "trade_volume_profile", "trade_arrival_rate",
|
||||
"compute_vpin", "compute_vpin_time_series", "fill_toxicity",
|
||||
"adverse_selection_ratio", "liquidation_clustering",
|
||||
"funding_regime", "funding_predictability", "funding_carry_pnl",
|
||||
"basis_spread", "basis_convergence_speed",
|
||||
"composite_signal", "SignalPipeline", "detect_hft_regime",
|
||||
]
|
||||
@@ -0,0 +1,235 @@
|
||||
"""
|
||||
Order book microstructure analytics.
|
||||
|
||||
Functions operate on book snapshots (bids/asks dicts or DataFrames)
|
||||
and return time-series or summary stats.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import math
|
||||
from collections import deque
|
||||
from typing import Optional
|
||||
|
||||
import numpy as np
|
||||
|
||||
|
||||
# ── Microprice ───────────────────────────────────────────────
|
||||
|
||||
def microprice(
|
||||
bids: dict[float, float],
|
||||
asks: dict[float, float],
|
||||
weight_bid: float = 0.5,
|
||||
) -> float:
|
||||
"""Weighted mid based on depth imbalance.
|
||||
|
||||
microprice = weight * bid_side + (1-weight) * ask_side
|
||||
where weight = bid_depth / (bid_depth + ask_depth)
|
||||
|
||||
Falls back to simple mid if no depth.
|
||||
"""
|
||||
bid_prices = sorted(bids.keys(), reverse=True)
|
||||
ask_prices = sorted(asks.keys())
|
||||
|
||||
if not bid_prices or not ask_prices:
|
||||
return 0.0
|
||||
|
||||
best_bid = bid_prices[0]
|
||||
best_ask = ask_prices[0]
|
||||
|
||||
bid_vol = sum(bids[px] for px in bid_prices[:10])
|
||||
ask_vol = sum(asks[px] for px in ask_prices[:10])
|
||||
total = bid_vol + ask_vol
|
||||
if total == 0:
|
||||
return (best_bid + best_ask) / 2.0
|
||||
|
||||
w = bid_vol / total
|
||||
return w * best_bid + (1 - w) * best_ask
|
||||
|
||||
|
||||
def mid_price(bids: dict[float, float], asks: dict[float, float]) -> float:
|
||||
bid_prices = sorted(bids.keys(), reverse=True)
|
||||
ask_prices = sorted(asks.keys())
|
||||
if not bid_prices or not ask_prices:
|
||||
return 0.0
|
||||
return (bid_prices[0] + ask_prices[0]) / 2.0
|
||||
|
||||
|
||||
# ── Order-book imbalance ─────────────────────────────────────
|
||||
|
||||
def order_book_imbalance(
|
||||
bids: dict[float, float],
|
||||
asks: dict[float, float],
|
||||
levels: int = 10,
|
||||
) -> float:
|
||||
"""OBI = (bid_vol - ask_vol) / (bid_vol + ask_vol). Range [-1, 1]."""
|
||||
bid_prices = sorted(bids.keys(), reverse=True)[:levels]
|
||||
ask_prices = sorted(asks.keys())[:levels]
|
||||
bid_vol = sum(bids[px] for px in bid_prices)
|
||||
ask_vol = sum(asks[px] for px in ask_prices)
|
||||
total = bid_vol + ask_vol
|
||||
if total == 0:
|
||||
return 0.0
|
||||
return (bid_vol - ask_vol) / total
|
||||
|
||||
|
||||
def depth_imbalance(
|
||||
bids: dict[float, float],
|
||||
asks: dict[float, float],
|
||||
price_distance_pct: float = 0.01,
|
||||
) -> float:
|
||||
"""Imbalance at a fixed price distance from mid (percentage-based)."""
|
||||
mid = mid_price(bids, asks)
|
||||
if mid <= 0:
|
||||
return 0.0
|
||||
lo = mid * (1 - price_distance_pct)
|
||||
hi = mid * (1 + price_distance_pct)
|
||||
bid_vol = sum(sz for px, sz in bids.items() if px >= lo)
|
||||
ask_vol = sum(sz for px, sz in asks.items() if px <= hi)
|
||||
total = bid_vol + ask_vol
|
||||
if total == 0:
|
||||
return 0.0
|
||||
return (bid_vol - ask_vol) / total
|
||||
|
||||
|
||||
# ── Spread statistics ────────────────────────────────────────
|
||||
|
||||
def spread_stats(
|
||||
bids: dict[float, float],
|
||||
asks: dict[float, float],
|
||||
) -> dict:
|
||||
bid_prices = sorted(bids.keys(), reverse=True)
|
||||
ask_prices = sorted(asks.keys())
|
||||
if not bid_prices or not ask_prices:
|
||||
return {"spread": 0, "spread_bps": 0, "mid": 0, "bid": 0, "ask": 0}
|
||||
|
||||
best_bid = bid_prices[0]
|
||||
best_ask = ask_prices[0]
|
||||
mid = (best_bid + best_ask) / 2.0
|
||||
spread = best_ask - best_bid
|
||||
spread_bps = (spread / mid * 10000) if mid > 0 else 0
|
||||
return {
|
||||
"spread": round(spread, 2),
|
||||
"spread_bps": round(spread_bps, 2),
|
||||
"mid": round(mid, 2),
|
||||
"best_bid": best_bid,
|
||||
"best_ask": best_ask,
|
||||
}
|
||||
|
||||
|
||||
# ── Depth resiliency ──────────────────────────────────────────
|
||||
|
||||
def depth_resiliency(
|
||||
bids: dict[float, float],
|
||||
asks: dict[float, float],
|
||||
impact_bps: float = 10.0,
|
||||
) -> dict:
|
||||
"""How much size sits within N bps of mid on each side.
|
||||
|
||||
Returns volume and level count within the impact radius, plus
|
||||
a resiliency score: bid_depth / ask_depth (ratio).
|
||||
"""
|
||||
mid = mid_price(bids, asks)
|
||||
if mid <= 0:
|
||||
return {"bid_vol": 0, "ask_vol": 0, "bid_levels": 0, "ask_levels": 0, "resiliency": 0}
|
||||
|
||||
radius = mid * impact_bps / 10000
|
||||
bid_lo = mid - radius
|
||||
ask_hi = mid + radius
|
||||
|
||||
bid_vol = sum(sz for px, sz in bids.items() if px >= bid_lo)
|
||||
ask_vol = sum(sz for px, sz in asks.items() if px <= ask_hi)
|
||||
bid_levels = sum(1 for px in bids if px >= bid_lo)
|
||||
ask_levels = sum(1 for px in asks if px <= ask_hi)
|
||||
|
||||
resiliency = bid_vol / ask_vol if ask_vol > 0 else float("inf")
|
||||
return {
|
||||
"bid_vol": round(bid_vol, 8),
|
||||
"ask_vol": round(ask_vol, 8),
|
||||
"bid_levels": bid_levels,
|
||||
"ask_levels": ask_levels,
|
||||
"resiliency": round(resiliency, 4),
|
||||
}
|
||||
|
||||
|
||||
def queue_depletion_prob(
|
||||
bids: dict[float, float],
|
||||
asks: dict[float, float],
|
||||
level_distance: int = 0,
|
||||
trade_rate_per_sec: float = 1.0,
|
||||
avg_trade_size: float = 0.01,
|
||||
) -> float:
|
||||
"""Probability queued order at best level(s) gets filled within 1 second.
|
||||
|
||||
Simple Poisson model: P(fill) = 1 - exp(-λ)
|
||||
where λ = trade_rate * avg_trade_size / depth_at_level
|
||||
|
||||
level_distance=0 means best bid/ask, 1 = one level behind, etc.
|
||||
"""
|
||||
bid_prices = sorted(bids.keys(), reverse=True)
|
||||
ask_prices = sorted(asks.keys())
|
||||
if not bid_prices or not ask_prices:
|
||||
return 0.0
|
||||
|
||||
px = bid_prices[min(level_distance, len(bid_prices) - 1)]
|
||||
depth = bids.get(px, 0)
|
||||
if depth <= 0:
|
||||
return 0.0
|
||||
|
||||
lam = trade_rate_per_sec * avg_trade_size / depth
|
||||
prob = 1.0 - math.exp(-lam)
|
||||
return round(min(prob, 0.9999), 6)
|
||||
|
||||
|
||||
# ── Batch processing ─────────────────────────────────────────
|
||||
|
||||
def batch_book_stats(
|
||||
snapshots: list[dict],
|
||||
levels: int = 10,
|
||||
) -> dict:
|
||||
"""Process a list of book snapshots and return aggregate stats.
|
||||
|
||||
Each snapshot: {"bids": {price: size, ...}, "asks": {price: size, ...}}
|
||||
"""
|
||||
obis = []
|
||||
spreads_bps = []
|
||||
micros = []
|
||||
depths = []
|
||||
resiliencies = []
|
||||
|
||||
for snap in snapshots:
|
||||
bids = snap.get("bids", {})
|
||||
asks = snap.get("asks", {})
|
||||
if not bids or not asks:
|
||||
continue
|
||||
obis.append(order_book_imbalance(bids, asks, levels))
|
||||
ss = spread_stats(bids, asks)
|
||||
spreads_bps.append(ss["spread_bps"])
|
||||
micros.append(microprice(bids, asks))
|
||||
dr = depth_resiliency(bids, asks)
|
||||
depths.append(dr["bid_vol"] + dr["ask_vol"])
|
||||
resiliencies.append(dr["resiliency"])
|
||||
|
||||
def _summarize(vals):
|
||||
if not vals:
|
||||
return {"mean": 0, "median": 0, "std": 0, "min": 0, "max": 0, "count": 0}
|
||||
a = np.array(vals, dtype=float)
|
||||
a = a[np.isfinite(a)]
|
||||
return {
|
||||
"mean": round(float(np.mean(a)), 4),
|
||||
"median": round(float(np.median(a)), 4),
|
||||
"std": round(float(np.std(a)), 4),
|
||||
"min": round(float(np.min(a)), 4),
|
||||
"max": round(float(np.max(a)), 4),
|
||||
"count": len(a),
|
||||
}
|
||||
|
||||
return {
|
||||
"obi": _summarize(obis),
|
||||
"spread_bps": _summarize(spreads_bps),
|
||||
"microprice_ratio": _summarize([m / mid_price(s["bids"], s["asks"])
|
||||
if mid_price(s["bids"], s["asks"]) > 0 else 1.0
|
||||
for s, m in zip(snapshots, micros)]),
|
||||
"depth_total": _summarize(depths),
|
||||
"resiliency": _summarize([r for r in resiliencies if r < 1e6]),
|
||||
}
|
||||
@@ -0,0 +1,179 @@
|
||||
"""
|
||||
Funding and basis behavior analytics.
|
||||
|
||||
Analyses funding rate regimes, basis dynamics (spot vs perp),
|
||||
and carry trade profitability.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import math
|
||||
|
||||
import numpy as np
|
||||
|
||||
|
||||
# ── Funding rate analytics ──────────────────────────────────
|
||||
|
||||
def funding_regime(
|
||||
funding_rates: list[float],
|
||||
window_hours: int = 24,
|
||||
n_samples_per_hour: int = 60, # e.g., 1 sample/min → 60/hr
|
||||
) -> dict:
|
||||
"""Classify current funding regime.
|
||||
|
||||
Returns regime classification and rolling stats.
|
||||
"""
|
||||
window = window_hours * n_samples_per_hour
|
||||
if len(funding_rates) < window:
|
||||
return {"regime": "insufficient_data", "mean_annual": 0, "volatility": 0}
|
||||
|
||||
recent = funding_rates[-window:]
|
||||
mean_rate = float(np.mean(recent))
|
||||
std_rate = float(np.std(recent))
|
||||
|
||||
# Funding is per-hour rate. Annualize: compounded 3x daily.
|
||||
# Hyperliquid funding: 8h rate * 3 = daily, * 365 = annual (approx)
|
||||
ann_rate = mean_rate * 3 * 365 * 100 # *100 to convert from fraction to %
|
||||
|
||||
if ann_rate > 15:
|
||||
regime = "high_positive"
|
||||
elif ann_rate > 5:
|
||||
regime = "positive"
|
||||
elif ann_rate < -15:
|
||||
regime = "high_negative"
|
||||
elif ann_rate < -5:
|
||||
regime = "negative"
|
||||
else:
|
||||
regime = "neutral"
|
||||
|
||||
return {
|
||||
"regime": regime,
|
||||
"mean_hourly": round(float(mean_rate), 8),
|
||||
"mean_annual_pct": round(ann_rate, 2),
|
||||
"volatility": round(float(std_rate), 8),
|
||||
"window_hours": window_hours,
|
||||
}
|
||||
|
||||
|
||||
def funding_predictability(
|
||||
funding_history: list[float],
|
||||
n_lags: int = 3,
|
||||
) -> dict:
|
||||
"""Measure funding rate autocorrelation — is funding momentum persistent?"""
|
||||
if len(funding_history) < n_lags + 2:
|
||||
return {"autocorr": [], "momentum_strength": 0}
|
||||
|
||||
autocorr = []
|
||||
for lag in range(1, n_lags + 1):
|
||||
x = funding_history[:-lag]
|
||||
y = funding_history[lag:]
|
||||
if len(x) < 2:
|
||||
autocorr.append(0)
|
||||
continue
|
||||
corr = np.corrcoef(x, y)[0, 1]
|
||||
autocorr.append(round(float(corr) if not np.isnan(corr) else 0, 4))
|
||||
|
||||
momentum = float(np.mean([abs(a) for a in autocorr]))
|
||||
|
||||
return {
|
||||
"autocorr": autocorr,
|
||||
"momentum_strength": round(momentum, 4),
|
||||
"is_momentum": momentum > 0.3,
|
||||
}
|
||||
|
||||
|
||||
def funding_carry_pnl(
|
||||
funding_rates: list[float],
|
||||
position_size: float = 1.0,
|
||||
mark_prices: list[float] | None = None,
|
||||
n_samples_per_hour: int = 60,
|
||||
) -> dict:
|
||||
"""Estimate carry PnL from holding a position given funding rates.
|
||||
|
||||
For positive funding → shorts earn, longs pay.
|
||||
"""
|
||||
if not funding_rates:
|
||||
return {"cumulative_pnl": 0, "hourly_pnl": []}
|
||||
|
||||
hourly = []
|
||||
cum = 0.0
|
||||
for i, rate in enumerate(funding_rates):
|
||||
if i % n_samples_per_hour == 0:
|
||||
notional = position_size * (mark_prices[i] if mark_prices and i < len(mark_prices) else 1.0)
|
||||
pnl = rate * notional # funding rate * position notional
|
||||
cum += pnl
|
||||
hourly.append(round(pnl, 8))
|
||||
|
||||
return {
|
||||
"cumulative_pnl": round(cum, 6),
|
||||
"hourly_pnl": hourly[-72:], # last 72 hours
|
||||
"n_hours": len(hourly),
|
||||
}
|
||||
|
||||
|
||||
# ── Basis analytics ──────────────────────────────────────────
|
||||
|
||||
def basis_spread(
|
||||
perp_prices: list[float],
|
||||
spot_prices: list[float],
|
||||
) -> dict:
|
||||
"""Compute basis (perp premium over spot) and its statistics.
|
||||
|
||||
basis_bps = (perp - spot) / spot * 10000
|
||||
"""
|
||||
min_len = min(len(perp_prices), len(spot_prices))
|
||||
if min_len < 2:
|
||||
return {"current_basis_bps": 0, "mean_basis_bps": 0, "max_basis_bps": 0}
|
||||
|
||||
perp = perp_prices[-min_len:]
|
||||
spot = spot_prices[-min_len:]
|
||||
basis_arr = []
|
||||
|
||||
for p, s in zip(perp, spot):
|
||||
if s > 0:
|
||||
basis_arr.append((p - s) / s * 10000)
|
||||
|
||||
a = np.array(basis_arr) if basis_arr else np.array([0.0])
|
||||
|
||||
return {
|
||||
"current_basis_bps": round(float(a[-1]), 2) if len(a) > 0 else 0,
|
||||
"mean_basis_bps": round(float(np.mean(a)), 2),
|
||||
"std_basis_bps": round(float(np.std(a)), 2),
|
||||
"max_basis_bps": round(float(np.max(a)), 2),
|
||||
"min_basis_bps": round(float(np.min(a)), 2),
|
||||
"n_samples": len(a),
|
||||
}
|
||||
|
||||
|
||||
def basis_convergence_speed(
|
||||
basis_history: list[float],
|
||||
half_life_lookback: int = 1440, # 24h at 1min samples
|
||||
) -> dict:
|
||||
"""Estimate basis mean-reversion half-life via AR(1).
|
||||
|
||||
Half-life = -log(2) / log(|rho|)
|
||||
"""
|
||||
if len(basis_history) < 10:
|
||||
return {"half_life_minutes": 0, "ar1_coef": 0, "mean_reverting": False}
|
||||
|
||||
x = basis_history[-half_life_lookback:]
|
||||
if len(x) < 10:
|
||||
return {"half_life_minutes": 0, "ar1_coef": 0, "mean_reverting": False}
|
||||
|
||||
x_t = x[:-1]
|
||||
x_t1 = x[1:]
|
||||
|
||||
rho = np.corrcoef(x_t, x_t1)[0, 1]
|
||||
rho = max(min(rho, 0.999), -0.999)
|
||||
|
||||
if abs(rho) < 0.01:
|
||||
half_life = 0
|
||||
else:
|
||||
half_life = -math.log(2) / math.log(abs(rho))
|
||||
|
||||
return {
|
||||
"half_life_minutes": round(half_life, 1),
|
||||
"ar1_coef": round(float(rho), 4),
|
||||
"mean_reverting": abs(rho) > 0.1 and rho < 0.95,
|
||||
}
|
||||
|
||||
@@ -0,0 +1,178 @@
|
||||
"""
|
||||
Composite signal construction from microstructure features.
|
||||
|
||||
Combines book imbalance, trade flow, toxicity, and funding signals
|
||||
into a single directional signal with confidence score.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Optional
|
||||
|
||||
import numpy as np
|
||||
|
||||
|
||||
# ── Composite signal ────────────────────────────────────────
|
||||
|
||||
def composite_signal(
|
||||
obi: float, # [-1, 1] order book imbalance
|
||||
trade_imbalance: float = 0.0, # [-1, 1] recent trade aggressor skew
|
||||
vpin: float = 0.0, # [0, 1] flow toxicity (higher = toxic)
|
||||
funding_regime: str = "neutral", # "high_positive", "negative", etc.
|
||||
spread_bps: float = 1.0, # current spread
|
||||
weight_obi: float = 0.35,
|
||||
weight_trade: float = 0.25,
|
||||
weight_vpin: float = -0.20, # negative: high VPIN → reduce confidence
|
||||
weight_funding: float = 0.20,
|
||||
) -> dict:
|
||||
"""Combine microstructure features into a single directional signal.
|
||||
|
||||
Returns:
|
||||
signal: "buy", "sell", or "neutral"
|
||||
score: [-1, 1] raw composite (positive = buy pressure)
|
||||
confidence: [0, 1] confidence in the signal
|
||||
breakdown: per-component contributions
|
||||
"""
|
||||
obi_score = np.clip(obi, -1.0, 1.0)
|
||||
trade_score = np.clip(trade_imbalance, -1.0, 1.0)
|
||||
vpin_score = np.clip(vpin, 0.0, 1.0)
|
||||
|
||||
# Funding: positive funding → short gets paid → sell bias; negative → buy bias
|
||||
funding_score_map = {
|
||||
"high_positive": -0.8,
|
||||
"positive": -0.4,
|
||||
"neutral": 0.0,
|
||||
"negative": 0.4,
|
||||
"high_negative": 0.8,
|
||||
}
|
||||
funding_score = funding_score_map.get(funding_regime, 0.0)
|
||||
|
||||
raw = (
|
||||
weight_obi * obi_score +
|
||||
weight_trade * trade_score +
|
||||
weight_vpin * vpin_score +
|
||||
weight_funding * funding_score
|
||||
)
|
||||
|
||||
# Confidence: base from signal magnitude, reduced by VPIN and spread
|
||||
signal_magnitude = abs(raw)
|
||||
vpin_penalty = np.clip(vpin * 0.5, 0.0, 0.3) if raw != 0 else 0
|
||||
spread_penalty = min(spread_bps / 50.0, 0.3) # wide spread → lower confidence
|
||||
|
||||
confidence = max(0.0, min(1.0, signal_magnitude * 1.5 - vpin_penalty - spread_penalty))
|
||||
|
||||
if raw > 0.1:
|
||||
signal = "buy"
|
||||
elif raw < -0.1:
|
||||
signal = "sell"
|
||||
else:
|
||||
signal = "neutral"
|
||||
|
||||
return {
|
||||
"signal": signal,
|
||||
"score": round(raw, 4),
|
||||
"confidence": round(confidence, 4),
|
||||
"breakdown": {
|
||||
"obi": round(obi_score * weight_obi, 4),
|
||||
"trade": round(trade_score * weight_trade, 4),
|
||||
"vpin": round(vpin_score * weight_vpin, 4),
|
||||
"funding": round(funding_score * weight_funding, 4),
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
# ── Signal pipeline ─────────────────────────────────────────
|
||||
|
||||
class SignalPipeline:
|
||||
"""Stateful pipeline that accumulates microstructure data and emits signals.
|
||||
|
||||
Usage:
|
||||
pipeline = SignalPipeline()
|
||||
pipeline.update_book(bids, asks)
|
||||
pipeline.update_trade(px, sz, mid)
|
||||
signal = pipeline.emit()
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
obi_window: int = 100,
|
||||
trade_window: int = 500,
|
||||
vpin_volume_size: float = 100.0,
|
||||
vpin_buckets: int = 50,
|
||||
):
|
||||
self._obi_window = obi_window
|
||||
self._trade_window = trade_window
|
||||
self._vpin_volume_size = vpin_volume_size
|
||||
self._vpin_buckets = vpin_buckets
|
||||
|
||||
self._buy_vol: list[float] = []
|
||||
self._sell_vol: list[float] = []
|
||||
self._obuys: int = 0
|
||||
self._osells: int = 0
|
||||
self._obuy: int = 1
|
||||
|
||||
def update_book(self, bids: dict[float, float], asks: dict[float, float]):
|
||||
from microstructure.book import order_book_imbalance
|
||||
self._obuys = order_book_imbalance(bids, asks)
|
||||
|
||||
def update_trade(self, px: float, sz: float, mid: float):
|
||||
if px >= mid:
|
||||
self._buy_vol.append(sz)
|
||||
else:
|
||||
self._sell_vol.append(sz)
|
||||
|
||||
if len(self._buy_vol) > self._trade_window:
|
||||
self._buy_vol = self._buy_vol[-self._trade_window:]
|
||||
if len(self._sell_vol) > self._trade_window:
|
||||
self._sell_vol = self._sell_vol[-self._trade_window:]
|
||||
|
||||
def recent_trade_imbalance(self) -> float:
|
||||
bv = sum(self._buy_vol[-100:])
|
||||
sv = sum(self._sell_vol[-100:])
|
||||
total = bv + sv
|
||||
return (bv - sv) / total if total > 0 else 0.0
|
||||
|
||||
def current_vpin(self) -> float:
|
||||
from microstructure.toxicity import compute_vpin
|
||||
result = compute_vpin(
|
||||
self._buy_vol, self._sell_vol,
|
||||
volume_bucket_size=self._vpin_volume_size,
|
||||
n_buckets=self._vpin_buckets,
|
||||
)
|
||||
return result.get("vpin_value", 0.0)
|
||||
|
||||
def emit(self) -> dict:
|
||||
return composite_signal(
|
||||
obi=self._obuys,
|
||||
trade_imbalance=self.recent_trade_imbalance(),
|
||||
vpin=self.current_vpin(),
|
||||
)
|
||||
|
||||
|
||||
# ── Regime detection ─────────────────────────────────────────
|
||||
|
||||
def detect_hft_regime(
|
||||
obi_std: float,
|
||||
spread_mean_bps: float,
|
||||
trade_rate_per_sec: float,
|
||||
vpin: float,
|
||||
) -> str:
|
||||
"""Classify current market regime for HFT strategy selection.
|
||||
|
||||
Returns one of:
|
||||
- "trending" — directional, high OBI variance, low VPIN
|
||||
- "ranging" — low OBI variance, tight spread, active
|
||||
- "toxic" — high VPIN, wide spread → don't quote
|
||||
- "quiet" — low activity, avoid
|
||||
"""
|
||||
if vpin > 0.4:
|
||||
return "toxic"
|
||||
if trade_rate_per_sec < 0.1:
|
||||
return "quiet"
|
||||
if obi_std > 0.3 and spread_mean_bps < 5:
|
||||
return "trending"
|
||||
if obi_std < 0.15 and spread_mean_bps < 3:
|
||||
return "ranging"
|
||||
if spread_mean_bps > 10:
|
||||
return "toxic"
|
||||
return "quiet"
|
||||
@@ -0,0 +1,248 @@
|
||||
"""
|
||||
Flow toxicity and adverse selection analytics.
|
||||
|
||||
VPIN (Volume-synchronized Probability of Informed Trading),
|
||||
fill toxicity metrics, and adverse selection indicators based
|
||||
on order book and trade data.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import math
|
||||
from collections import deque
|
||||
|
||||
import numpy as np
|
||||
|
||||
|
||||
# ── VPIN ─────────────────────────────────────────────────────
|
||||
|
||||
def compute_vpin(
|
||||
buy_volume: list[float],
|
||||
sell_volume: list[float],
|
||||
volume_bucket_size: float | None = None,
|
||||
n_buckets: int = 50,
|
||||
) -> dict:
|
||||
"""Volume-synchronized Probability of Informed Trading.
|
||||
|
||||
Args:
|
||||
buy_volume: volume classified as buyer-initiated per bar/period
|
||||
sell_volume: volume classified as seller-initiated per bar/period
|
||||
volume_bucket_size: target volume per bucket (auto if None)
|
||||
n_buckets: number of buckets for rolling VPIN
|
||||
|
||||
Returns dict with vpin values and summary.
|
||||
"""
|
||||
if not buy_volume or len(buy_volume) != len(sell_volume):
|
||||
return {"vpin_value": 0, "n_buckets": 0, "bucket_size": 0}
|
||||
|
||||
if volume_bucket_size is None:
|
||||
total_vol = sum(buy_volume) + sum(sell_volume)
|
||||
volume_bucket_size = total_vol / max(len(buy_volume), 1) * 5
|
||||
|
||||
buy_sell = [(b, s) for b, s in zip(buy_volume, sell_volume)]
|
||||
buckets: list[dict] = []
|
||||
current_buy = 0.0
|
||||
current_sell = 0.0
|
||||
|
||||
for b, s in buy_sell:
|
||||
current_buy += b
|
||||
current_sell += s
|
||||
if current_buy + current_sell >= volume_bucket_size:
|
||||
total = current_buy + current_sell
|
||||
buckets.append({
|
||||
"buy": current_buy,
|
||||
"sell": current_sell,
|
||||
"total": total,
|
||||
"imbalance": abs(current_buy - current_sell),
|
||||
})
|
||||
# Carry over excess
|
||||
excess = total - volume_bucket_size
|
||||
current_buy = excess * (current_buy / total) if total > 0 else 0
|
||||
current_sell = excess * (current_sell / total) if total > 0 else 0
|
||||
|
||||
vpin_value = 0.0
|
||||
if len(buckets) >= n_buckets:
|
||||
recent = buckets[-n_buckets:]
|
||||
total_imb = sum(b["imbalance"] for b in recent)
|
||||
total_vol = sum(b["total"] for b in recent)
|
||||
vpin_value = total_imb / total_vol if total_vol > 0 else 0.0
|
||||
|
||||
return {
|
||||
"vpin_value": round(vpin_value, 4),
|
||||
"n_buckets": len(buckets),
|
||||
"bucket_size": round(volume_bucket_size, 2),
|
||||
}
|
||||
|
||||
|
||||
def compute_vpin_time_series(
|
||||
buy_volume: list[float],
|
||||
sell_volume: list[float],
|
||||
volume_bucket_size: float | None = None,
|
||||
n_buckets: int = 50,
|
||||
) -> dict:
|
||||
"""Compute rolling VPIN time series."""
|
||||
if not buy_volume or len(buy_volume) != len(sell_volume):
|
||||
return {"vpin_values": [], "mean": 0, "std": 0, "max": 0, "threshold_alarm": 0}
|
||||
|
||||
if volume_bucket_size is None:
|
||||
total_vol = sum(buy_volume) + sum(sell_volume)
|
||||
volume_bucket_size = total_vol / max(len(buy_volume), 1) * 5
|
||||
|
||||
buy_sell = [(b, s) for b, s in zip(buy_volume, sell_volume)]
|
||||
buckets: list[float] = []
|
||||
current_buy = 0.0
|
||||
current_sell = 0.0
|
||||
vpin_series = []
|
||||
|
||||
for b, s in buy_sell:
|
||||
current_buy += b
|
||||
current_sell += s
|
||||
if current_buy + current_sell >= volume_bucket_size:
|
||||
total = current_buy + current_sell
|
||||
imb = abs(current_buy - current_sell)
|
||||
buckets.append(imb / total if total > 0 else 0.5)
|
||||
|
||||
excess = total - volume_bucket_size
|
||||
ratio = current_buy / total if total > 0 else 0.5
|
||||
current_buy = ratio * excess
|
||||
current_sell = (1 - ratio) * excess
|
||||
|
||||
if len(buckets) >= n_buckets:
|
||||
vpin_series.append(sum(buckets[-n_buckets:]) / n_buckets)
|
||||
else:
|
||||
vpin_series.append(sum(buckets) / len(buckets))
|
||||
|
||||
a = np.array(vpin_series) if vpin_series else np.array([0.0])
|
||||
return {
|
||||
"vpin_values": [round(v, 4) for v in vpin_series],
|
||||
"mean": round(float(np.mean(a)), 4),
|
||||
"std": round(float(np.std(a)), 4),
|
||||
"max": round(float(np.max(a)), 4),
|
||||
"threshold_alarm": round(float(np.mean(a) + 2 * np.std(a)), 4),
|
||||
}
|
||||
|
||||
|
||||
# ── Fill toxicity ────────────────────────────────────────────
|
||||
|
||||
def fill_toxicity(
|
||||
trade_prices: list[float],
|
||||
mids: list[float],
|
||||
trade_sides: list[str],
|
||||
horizon_ticks: int = 10,
|
||||
) -> dict:
|
||||
"""Compute toxicity per trade: did price move against you after fill?
|
||||
|
||||
For each buy trade: toxicity = (mid before - mid after) / mid
|
||||
For each sell: toxicity = (mid after - mid before) / mid
|
||||
|
||||
Positive toxicity = adverse price movement post-trade.
|
||||
"""
|
||||
if len(trade_prices) < horizon_ticks + 1:
|
||||
return {"buy_toxicity_mean": 0, "sell_toxicity_mean": 0, "overall": 0}
|
||||
|
||||
buy_tox = []
|
||||
sell_tox = []
|
||||
n = len(trade_prices)
|
||||
|
||||
for i in range(n - horizon_ticks):
|
||||
side = trade_sides[i] if i < len(trade_sides) else "unknown"
|
||||
mid_before = mids[min(i + 1, n - 1)]
|
||||
mid_after = mids[min(i + horizon_ticks, n - 1)]
|
||||
if mid_before <= 0 or mid_after <= 0:
|
||||
continue
|
||||
|
||||
change = (mid_before - mid_after) / mid_before
|
||||
|
||||
if side == "buy":
|
||||
buy_tox.append(change)
|
||||
elif side == "sell":
|
||||
sell_tox.append(-change)
|
||||
|
||||
return {
|
||||
"buy_toxicity_mean_bps": round(float(np.mean(buy_tox)) * 10000, 2) if buy_tox else 0,
|
||||
"sell_toxicity_mean_bps": round(float(np.mean(sell_tox)) * 10000, 2) if sell_tox else 0,
|
||||
"buy_count": len(buy_tox),
|
||||
"sell_count": len(sell_tox),
|
||||
"overall_bps": round(float(np.mean(buy_tox + sell_tox)) * 10000, 2) if buy_tox or sell_tox else 0,
|
||||
}
|
||||
|
||||
|
||||
# ── Adverse selection ───────────────────────────────────────
|
||||
|
||||
def adverse_selection_ratio(
|
||||
mid_after_trades: list[float],
|
||||
mid_before_trades: list[float],
|
||||
trade_sides: list[str],
|
||||
) -> dict:
|
||||
"""Adverse selection ratio per side (mid after / mid before - 1).
|
||||
|
||||
Higher values = more adverse selection (price moves against you).
|
||||
"""
|
||||
buys = []
|
||||
sells = []
|
||||
for i, side in enumerate(trade_sides):
|
||||
if i >= len(mid_after_trades) or i >= len(mid_before_trades):
|
||||
break
|
||||
before = mid_before_trades[i]
|
||||
after = mid_after_trades[i]
|
||||
if before <= 0:
|
||||
continue
|
||||
sel = (after - before) / before
|
||||
if side == "buy":
|
||||
buys.append(-sel)
|
||||
elif side == "sell":
|
||||
sells.append(sel)
|
||||
|
||||
ba = np.array(buys) if buys else np.array([0.0])
|
||||
sa = np.array(sells) if sells else np.array([0.0])
|
||||
|
||||
return {
|
||||
"buy_adverse_bps": round(float(np.mean(ba)) * 10000, 2),
|
||||
"sell_adverse_bps": round(float(np.mean(sa)) * 10000, 2),
|
||||
"buy_win_pct": round(np.sum(ba <= 0) / len(ba), 4) if len(ba) > 0 else 0,
|
||||
"sell_win_pct": round(np.sum(sa <= 0) / len(sa), 4) if len(sa) > 0 else 0,
|
||||
}
|
||||
|
||||
|
||||
# ── Liquidation clustering ──────────────────────────────────
|
||||
|
||||
def liquidation_clustering(
|
||||
liquidation_times_ms: list[int],
|
||||
window_sec: int = 300,
|
||||
) -> dict:
|
||||
"""Detect liquidation clusters — unusual concentration of liquidations.
|
||||
|
||||
Returns cluster periods and intensity.
|
||||
"""
|
||||
if len(liquidation_times_ms) < 2:
|
||||
return {"clusters": [], "mean_interval_s": 0, "clustered_pct": 0}
|
||||
|
||||
intervals = [
|
||||
(liquidation_times_ms[i + 1] - liquidation_times_ms[i]) / 1000
|
||||
for i in range(len(liquidation_times_ms) - 1)
|
||||
]
|
||||
mean_interval = float(np.mean(intervals))
|
||||
std_interval = float(np.std(intervals))
|
||||
|
||||
clusters = []
|
||||
cluster_start = None
|
||||
for i, interval in enumerate(intervals):
|
||||
if interval < mean_interval * 0.3: # Tight clustering threshold
|
||||
if cluster_start is None:
|
||||
cluster_start = liquidation_times_ms[i]
|
||||
else:
|
||||
if cluster_start is not None and liquidation_times_ms[i] - cluster_start < window_sec * 1000:
|
||||
clusters.append({
|
||||
"start_ms": cluster_start,
|
||||
"end_ms": liquidation_times_ms[i],
|
||||
"count": i - liquidation_times_ms.index(cluster_start) + 1 if cluster_start in liquidation_times_ms else 0,
|
||||
})
|
||||
cluster_start = None
|
||||
|
||||
clustered_count = sum(c.get("count", 0) for c in clusters)
|
||||
return {
|
||||
"clusters": clusters,
|
||||
"n_clusters": len(clusters),
|
||||
"mean_interval_s": round(mean_interval, 2),
|
||||
"clustered_pct": round(clustered_count / len(liquidation_times_ms), 4) if liquidation_times_ms else 0,
|
||||
}
|
||||
@@ -0,0 +1,202 @@
|
||||
"""
|
||||
Trade microstructure analytics — aggressor classification and markout curves.
|
||||
|
||||
Lee-Ready algorithm for trade direction classification, plus
|
||||
forward markout analysis: what happens to mid price N seconds after
|
||||
a trade of a given type.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import numpy as np
|
||||
|
||||
|
||||
# ── Aggressor classification ─────────────────────────────────
|
||||
|
||||
def classify_lee_ready(
|
||||
trade_px: float,
|
||||
mid_at_trade: float,
|
||||
bid_at_trade: float | None = None,
|
||||
ask_at_trade: float | None = None,
|
||||
) -> str:
|
||||
"""Lee-Ready: trade above mid = buy, below mid = sell.
|
||||
At mid: compare to previous tick (quote rule) — if unavailable,
|
||||
compare to bid/ask (trade at bid = sell, at ask = buy).
|
||||
"""
|
||||
if trade_px > mid_at_trade:
|
||||
return "buy"
|
||||
elif trade_px < mid_at_trade:
|
||||
return "sell"
|
||||
else:
|
||||
if ask_at_trade is not None and trade_px >= ask_at_trade:
|
||||
return "buy"
|
||||
if bid_at_trade is not None and trade_px <= bid_at_trade:
|
||||
return "sell"
|
||||
return "unknown"
|
||||
|
||||
|
||||
def classify_bulk_lee_ready(
|
||||
trades: list[dict],
|
||||
mids: list[float] | None = None,
|
||||
bids: list[float] | None = None,
|
||||
asks: list[float] | None = None,
|
||||
) -> list[str]:
|
||||
"""Classify a list of trades using Lee-Ready.
|
||||
|
||||
trades: [{"px": float, ...}, ...]
|
||||
mids: optional list of mid prices at each trade time
|
||||
bids/asks: optional best bid/ask at each trade time
|
||||
"""
|
||||
results = []
|
||||
for i, trade in enumerate(trades):
|
||||
px = float(trade.get("px", 0))
|
||||
mid = float(mids[i]) if mids and i < len(mids) else px
|
||||
bid = float(bids[i]) if bids and i < len(bids) else None
|
||||
ask = float(asks[i]) if asks and i < len(asks) else None
|
||||
results.append(classify_lee_ready(px, mid, bid, ask))
|
||||
return results
|
||||
|
||||
|
||||
# ── Markout curves ───────────────────────────────────────────
|
||||
|
||||
def compute_markouts(
|
||||
trades: list[dict],
|
||||
mid_prices: list[float],
|
||||
trade_times: list[int], # ms since epoch
|
||||
horizons_ms: list[int] | None = None,
|
||||
) -> dict:
|
||||
"""For each trade, compute mid-price change at specified horizons.
|
||||
|
||||
Returns:
|
||||
{"buys": {horizon_ms: [markout_values...]}, "sells": {...}, ...}
|
||||
"""
|
||||
if horizons_ms is None:
|
||||
horizons_ms = [100, 500, 1000, 5000, 10000, 30000, 60000]
|
||||
|
||||
results: dict[str, dict[int, list[float]]] = {
|
||||
"buy": {h: [] for h in horizons_ms},
|
||||
"sell": {h: [] for h in horizons_ms},
|
||||
}
|
||||
|
||||
bids_at_trade = []
|
||||
asks_at_trade = []
|
||||
mids_at_trade = []
|
||||
|
||||
for i, (trade, mid) in enumerate(zip(trades, mid_prices)):
|
||||
mids_at_trade.append(mid)
|
||||
bids_at_trade.append(mid * 0.9995 if mid > 0 else 0)
|
||||
asks_at_trade.append(mid * 1.0005 if mid > 0 else 0)
|
||||
|
||||
sides = classify_bulk_lee_ready(trades, mids_at_trade, bids_at_trade, asks_at_trade)
|
||||
|
||||
for i, (trade, side, t0) in enumerate(zip(trades, sides, trade_times)):
|
||||
base_mid = mid_prices[i] if i < len(mid_prices) else 0
|
||||
if base_mid <= 0:
|
||||
continue
|
||||
|
||||
for horizon in horizons_ms:
|
||||
target_ts = t0 + horizon
|
||||
future_mid = base_mid
|
||||
|
||||
for j in range(i + 1, len(mid_prices)):
|
||||
if trade_times[j] >= target_ts:
|
||||
future_mid = mid_prices[j]
|
||||
break
|
||||
else:
|
||||
if len(mid_prices) > i + 1:
|
||||
future_mid = mid_prices[-1]
|
||||
|
||||
markout = (future_mid - base_mid) / base_mid * 10000 # bps
|
||||
if side in ("buy", "sell"):
|
||||
results[side][horizon].append(markout)
|
||||
|
||||
return results
|
||||
|
||||
|
||||
def markout_summary(
|
||||
markouts: dict[str, dict[int, list[float]]],
|
||||
) -> dict:
|
||||
"""Summarize markout curves with mean, std, t-stat."""
|
||||
summary = {}
|
||||
for side in ("buy", "sell"):
|
||||
summary[side] = {}
|
||||
for horizon, vals in markouts.get(side, {}).items():
|
||||
if not vals:
|
||||
summary[side][horizon] = {"mean": 0, "std": 0, "t_stat": 0, "count": 0}
|
||||
continue
|
||||
a = np.array(vals, dtype=float)
|
||||
a = a[np.isfinite(a)]
|
||||
mean = float(np.mean(a))
|
||||
std = float(np.std(a, ddof=1))
|
||||
t_stat = mean / std * np.sqrt(len(a)) if std > 0 else 0
|
||||
summary[side][horizon] = {
|
||||
"mean_bps": round(mean, 2),
|
||||
"std_bps": round(std, 2),
|
||||
"t_stat": round(t_stat, 3),
|
||||
"count": len(a),
|
||||
}
|
||||
return summary
|
||||
|
||||
|
||||
# ── Trade metrics ────────────────────────────────────────────
|
||||
|
||||
def trade_volume_profile(
|
||||
trades: list[dict],
|
||||
n_buckets: int = 20,
|
||||
) -> dict:
|
||||
"""Volume profile: trade count and volume by size bucket."""
|
||||
sizes = [float(t.get("sz", 0)) for t in trades if float(t.get("sz", 0)) > 0]
|
||||
if not sizes:
|
||||
return {"buckets": [], "counts": [], "volumes": []}
|
||||
|
||||
min_sz, max_sz = min(sizes), max(sizes)
|
||||
if min_sz == max_sz:
|
||||
buckets = [min_sz]
|
||||
else:
|
||||
buckets = np.linspace(min_sz, max_sz, n_buckets + 1).tolist()
|
||||
|
||||
counts = [0] * n_buckets
|
||||
volumes = [0.0] * n_buckets
|
||||
for sz in sizes:
|
||||
for b in range(n_buckets):
|
||||
if buckets[b] <= sz < buckets[b + 1] or (b == n_buckets - 1 and sz == buckets[b + 1]):
|
||||
counts[b] += 1
|
||||
volumes[b] += sz
|
||||
break
|
||||
|
||||
return {
|
||||
"buckets": [round((buckets[i] + buckets[i + 1]) / 2, 6) for i in range(n_buckets)],
|
||||
"counts": counts,
|
||||
"volumes": [round(v, 6) for v in volumes],
|
||||
}
|
||||
|
||||
|
||||
def trade_arrival_rate(
|
||||
trade_times_ms: list[int],
|
||||
window_sec: int = 60,
|
||||
) -> dict:
|
||||
"""Trade arrival intensity (trades per second) over rolling windows."""
|
||||
if not trade_times_ms:
|
||||
return {"mean_rate": 0, "max_rate": 0, "burst_count": 0, "rates": []}
|
||||
|
||||
t0 = trade_times_ms[0]
|
||||
rates = []
|
||||
burst_count = 0
|
||||
window_ms = window_sec * 1000
|
||||
|
||||
for start in range(t0, trade_times_ms[-1], window_ms):
|
||||
end = start + window_ms
|
||||
count = sum(1 for t in trade_times_ms if start <= t < end)
|
||||
rate = count / window_sec
|
||||
rates.append(rate)
|
||||
if rate > rates[-2] * 3 if len(rates) > 1 else rate > 10:
|
||||
burst_count += 1
|
||||
|
||||
a = np.array(rates, dtype=float) if rates else np.array([0.0])
|
||||
return {
|
||||
"mean_rate": round(float(np.mean(a)), 3),
|
||||
"max_rate": round(float(np.max(a)), 3),
|
||||
"std_rate": round(float(np.std(a)), 3),
|
||||
"burst_count": burst_count,
|
||||
"rates": [round(r, 3) for r in rates[-100:]],
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
"""
|
||||
Tests for data/latency.py — latency tracking.
|
||||
"""
|
||||
import time
|
||||
|
||||
from data.latency import LatencyTracker
|
||||
|
||||
|
||||
def test_empty_stats():
|
||||
lt = LatencyTracker()
|
||||
s = lt.stats()
|
||||
for key in ("transport_ms", "signal_ms", "order_ms", "roundtrip_ms"):
|
||||
assert s[key]["count"] == 0
|
||||
|
||||
|
||||
def test_record_transport():
|
||||
lt = LatencyTracker()
|
||||
lt.record_transport(1705312800000, 1705312800.100) # 100ms
|
||||
s = lt.stats()
|
||||
assert s["transport_ms"]["p50"] == 100.0
|
||||
assert s["transport_ms"]["count"] == 1
|
||||
|
||||
|
||||
def test_record_signal():
|
||||
lt = LatencyTracker()
|
||||
lt.record_signal(5.2)
|
||||
lt.record_signal(3.1)
|
||||
s = lt.stats()
|
||||
assert s["signal_ms"]["p50"] == 4.15 # interpolated median of [3.1, 5.2]
|
||||
assert s["signal_ms"]["count"] == 2
|
||||
|
||||
|
||||
def test_record_order():
|
||||
lt = LatencyTracker()
|
||||
lt.record_order(12.5)
|
||||
s = lt.stats()
|
||||
assert s["order_ms"]["p50"] == 12.5
|
||||
|
||||
|
||||
def test_record_roundtrip():
|
||||
lt = LatencyTracker()
|
||||
lt.record_roundtrip(150.0)
|
||||
s = lt.stats()
|
||||
assert s["roundtrip_ms"]["p50"] == 150.0
|
||||
|
||||
|
||||
def test_negative_latency_ignored():
|
||||
lt = LatencyTracker()
|
||||
lt.record_transport(1705312800100, 1705312800.000) # local before exchange
|
||||
s = lt.stats()
|
||||
assert s["transport_ms"]["count"] == 0
|
||||
|
||||
|
||||
def test_summary_compact():
|
||||
lt = LatencyTracker()
|
||||
lt.record_transport(1705312800000, 1705312800.100)
|
||||
lt.record_signal(5.0)
|
||||
summary = lt.summary()
|
||||
assert "transport_ms" in summary
|
||||
assert "signal_ms" in summary
|
||||
assert "p50" in summary["transport_ms"]
|
||||
|
||||
|
||||
def test_percentiles_multiple_values():
|
||||
lt = LatencyTracker()
|
||||
for lat in [10, 20, 30, 40, 50, 60, 70, 80, 90, 100]:
|
||||
lt.record_signal(float(lat))
|
||||
s = lt.stats()
|
||||
assert s["signal_ms"]["p50"] == 55.0 # interpolated median of 10 values
|
||||
assert s["signal_ms"]["p90"] == 91.0 # (10-1)*90/100=8.1, lo=8, vals[8]=90, vals[9]=100, frac=0.1, result=91
|
||||
|
||||
|
||||
def test_window_pruning_approx():
|
||||
"""Old records should be pruned after window expires."""
|
||||
lt = LatencyTracker(window_seconds=0.1)
|
||||
lt.record_signal(10.0)
|
||||
time.sleep(0.15)
|
||||
lt.record_signal(20.0)
|
||||
s = lt.stats()
|
||||
assert s["signal_ms"]["p50"] == 20.0 # only the recent one
|
||||
@@ -0,0 +1,388 @@
|
||||
"""
|
||||
Tests for microstructure module — book, trades, toxicity, funding, signals.
|
||||
"""
|
||||
import math
|
||||
|
||||
import numpy as np
|
||||
|
||||
from microstructure.book import (
|
||||
microprice,
|
||||
mid_price,
|
||||
order_book_imbalance,
|
||||
depth_imbalance,
|
||||
spread_stats,
|
||||
depth_resiliency,
|
||||
queue_depletion_prob,
|
||||
batch_book_stats,
|
||||
)
|
||||
from microstructure.trades import (
|
||||
classify_lee_ready,
|
||||
classify_bulk_lee_ready,
|
||||
compute_markouts,
|
||||
markout_summary,
|
||||
trade_volume_profile,
|
||||
trade_arrival_rate,
|
||||
)
|
||||
from microstructure.toxicity import (
|
||||
compute_vpin,
|
||||
compute_vpin_time_series,
|
||||
fill_toxicity,
|
||||
adverse_selection_ratio,
|
||||
liquidation_clustering,
|
||||
)
|
||||
from microstructure.funding import (
|
||||
funding_regime,
|
||||
funding_predictability,
|
||||
funding_carry_pnl,
|
||||
basis_spread,
|
||||
basis_convergence_speed,
|
||||
)
|
||||
from microstructure.signals import (
|
||||
composite_signal,
|
||||
SignalPipeline,
|
||||
detect_hft_regime,
|
||||
)
|
||||
|
||||
|
||||
# ── Helpers ────────────────────────────────────────────────
|
||||
|
||||
def _basic_book():
|
||||
"""BTC-style book."""
|
||||
bids = {50000.0 - i * 10.0: 1.0 + i * 0.2 for i in range(20)}
|
||||
asks = {50001.0 + i * 10.0: 1.0 + i * 0.2 for i in range(20)}
|
||||
return bids, asks
|
||||
|
||||
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
# Book tests
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
|
||||
class TestMidPrice:
|
||||
def test_simple_mid(self):
|
||||
bids = {100.0: 1.0}
|
||||
asks = {102.0: 1.0}
|
||||
assert mid_price(bids, asks) == 101.0
|
||||
|
||||
def test_empty_returns_zero(self):
|
||||
assert mid_price({}, {}) == 0.0
|
||||
|
||||
|
||||
class TestMicroprice:
|
||||
def test_equal_depth(self):
|
||||
bids = {100.0: 1.0}
|
||||
asks = {102.0: 1.0}
|
||||
mid = microprice(bids, asks)
|
||||
assert mid == 101.0 # equal weight = simple mid
|
||||
|
||||
def test_bid_heavy(self):
|
||||
bids = {100.0: 10.0}
|
||||
asks = {102.0: 1.0}
|
||||
mid = microprice(bids, asks)
|
||||
assert mid < 101.0 # titled toward heavier bid side → lower price
|
||||
|
||||
def test_ask_heavy(self):
|
||||
bids = {100.0: 1.0}
|
||||
asks = {102.0: 10.0}
|
||||
mid = microprice(bids, asks)
|
||||
assert mid > 101.0 # tilted toward heavier ask side → higher price
|
||||
|
||||
|
||||
class TestOrderBookImbalance:
|
||||
def test_balanced(self):
|
||||
bids = {100.0: 5.0}
|
||||
asks = {102.0: 5.0}
|
||||
obi = order_book_imbalance(bids, asks)
|
||||
assert obi == 0.0
|
||||
|
||||
def test_bid_heavy(self):
|
||||
bids = {100.0: 8.0}
|
||||
asks = {102.0: 2.0}
|
||||
obi = order_book_imbalance(bids, asks)
|
||||
assert obi > 0
|
||||
assert obi == 0.6
|
||||
|
||||
def test_ask_heavy(self):
|
||||
bids = {100.0: 2.0}
|
||||
asks = {102.0: 8.0}
|
||||
obi = order_book_imbalance(bids, asks)
|
||||
assert obi < 0
|
||||
assert obi == -0.6
|
||||
|
||||
def test_empty_book(self):
|
||||
assert order_book_imbalance({}, {}) == 0.0
|
||||
|
||||
|
||||
class TestSpreadStats:
|
||||
def test_basic(self):
|
||||
bids = {50000.0: 1.0}
|
||||
asks = {50002.0: 1.0}
|
||||
stats = spread_stats(bids, asks)
|
||||
assert stats["spread"] == 2.0
|
||||
assert stats["best_bid"] == 50000.0
|
||||
assert stats["best_ask"] == 50002.0
|
||||
assert stats["spread_bps"] > 0
|
||||
|
||||
def test_empty(self):
|
||||
stats = spread_stats({}, {})
|
||||
assert stats["spread"] == 0
|
||||
|
||||
|
||||
class TestDepthResiliency:
|
||||
def test_basic(self):
|
||||
bids, asks = _basic_book()
|
||||
dr = depth_resiliency(bids, asks, impact_bps=100.0)
|
||||
assert dr["bid_vol"] > 0
|
||||
assert dr["ask_vol"] > 0
|
||||
assert dr["bid_levels"] > 0
|
||||
|
||||
|
||||
class TestQueueDepletionProb:
|
||||
def test_probability_range(self):
|
||||
bids = {50000.0: 1.0}
|
||||
asks = {50002.0: 0.5}
|
||||
prob = queue_depletion_prob(bids, asks, level_distance=0, trade_rate_per_sec=2.0, avg_trade_size=0.1)
|
||||
assert 0.0 <= prob <= 1.0
|
||||
|
||||
def test_deep_book_low_prob(self):
|
||||
bids = {50000.0: 100.0}
|
||||
asks = {50002.0: 100.0}
|
||||
prob = queue_depletion_prob(bids, asks, level_distance=0, trade_rate_per_sec=1.0, avg_trade_size=0.01)
|
||||
assert prob < 0.01
|
||||
|
||||
|
||||
class TestBatchBookStats:
|
||||
def test_multiple_snapshots(self):
|
||||
bids, asks = _basic_book()
|
||||
snaps = [{"bids": bids, "asks": asks} for _ in range(5)]
|
||||
result = batch_book_stats(snaps)
|
||||
assert result["obi"]["count"] == 5
|
||||
assert result["spread_bps"]["count"] == 5
|
||||
|
||||
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
# Trade tests
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
|
||||
class TestLeeReady:
|
||||
def test_buy_above_mid(self):
|
||||
assert classify_lee_ready(105.0, 100.0) == "buy"
|
||||
|
||||
def test_sell_below_mid(self):
|
||||
assert classify_lee_ready(95.0, 100.0) == "sell"
|
||||
|
||||
def test_at_mid_with_ask(self):
|
||||
assert classify_lee_ready(100.0, 100.0, bid_at_trade=99.0, ask_at_trade=100.0) == "buy"
|
||||
|
||||
def test_at_mid_with_bid(self):
|
||||
assert classify_lee_ready(100.0, 100.0, bid_at_trade=100.0, ask_at_trade=101.0) == "sell"
|
||||
|
||||
def test_unknown_at_mid(self):
|
||||
assert classify_lee_ready(100.0, 100.0) == "unknown"
|
||||
|
||||
|
||||
class TestBulkLeeReady:
|
||||
def test_bulk_classification(self):
|
||||
trades = [{"px": 105}, {"px": 95}, {"px": 100}]
|
||||
mids = [100.0, 100.0, 100.0]
|
||||
bids = [99.0, 99.0, 100.0]
|
||||
asks = [101.0, 101.0, 101.0]
|
||||
sides = classify_bulk_lee_ready(trades, mids, bids, asks)
|
||||
assert sides == ["buy", "sell", "sell"]
|
||||
|
||||
|
||||
class TestMarkouts:
|
||||
def test_basic_markout(self):
|
||||
trades = [{"px": 100.0}, {"px": 101.0}]
|
||||
mids = [100.0, 100.0, 100.1, 100.2]
|
||||
times = [0, 100, 200, 300]
|
||||
result = compute_markouts(trades, mids, times, horizons_ms=[100, 200])
|
||||
assert "buy" in result
|
||||
assert "sell" in result
|
||||
|
||||
def test_markout_summary(self):
|
||||
markouts = {"buy": {100: [1.0, 2.0, -1.0]}, "sell": {100: [-1.0, -2.0]}}
|
||||
summary = markout_summary(markouts)
|
||||
assert summary["buy"][100]["count"] == 3
|
||||
assert summary["sell"][100]["count"] == 2
|
||||
|
||||
|
||||
class TestTradeVolumeProfile:
|
||||
def test_basic(self):
|
||||
trades = [{"sz": 0.1}, {"sz": 0.2}, {"sz": 0.3}]
|
||||
profile = trade_volume_profile(trades, n_buckets=3)
|
||||
assert profile["buckets"]
|
||||
assert sum(profile["counts"]) == 3
|
||||
|
||||
|
||||
class TestTradeArrivalRate:
|
||||
def test_basic(self):
|
||||
times = list(range(0, 60000, 1000)) # 1 trade/sec for 60 sec
|
||||
result = trade_arrival_rate(times, window_sec=10)
|
||||
assert result["mean_rate"] > 0
|
||||
|
||||
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
# Toxicity tests
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
|
||||
class TestVPIN:
|
||||
def test_balanced_volume(self):
|
||||
buy_vol = [1.0] * 100
|
||||
sell_vol = [1.0] * 100
|
||||
result = compute_vpin(buy_vol, sell_vol)
|
||||
assert result["vpin_value"] >= 0
|
||||
|
||||
def test_imbalanced_volume(self):
|
||||
buy_vol = [2.0] * 200
|
||||
sell_vol = [1.0] * 200
|
||||
result = compute_vpin(buy_vol, sell_vol, volume_bucket_size=5.0, n_buckets=20)
|
||||
assert result["vpin_value"] > 0 # buy > sell → imbalance > 0
|
||||
|
||||
def test_empty_returns_zero(self):
|
||||
result = compute_vpin([], [])
|
||||
assert result["vpin_value"] == 0
|
||||
|
||||
|
||||
class TestVPINTimeSeries:
|
||||
def test_returns_series(self):
|
||||
buy_vol = [1.0] * 200 + [3.0] * 50
|
||||
sell_vol = [1.0] * 200 + [0.5] * 50
|
||||
result = compute_vpin_time_series(buy_vol, sell_vol, n_buckets=20)
|
||||
assert result["vpin_values"]
|
||||
assert result["mean"] >= 0
|
||||
assert result["max"] >= 0
|
||||
|
||||
|
||||
class TestFillToxicity:
|
||||
def test_basic_toxicity(self):
|
||||
prices = [100.0, 101.0, 102.0, 103.0, 104.0, 105.0]
|
||||
mids = [100.0, 100.5, 101.0, 101.5, 102.0, 102.5]
|
||||
sides = ["buy", "buy", "sell", "sell", "buy", "sell"]
|
||||
result = fill_toxicity(prices, mids, sides, horizon_ticks=2)
|
||||
assert "overall_bps" in result
|
||||
|
||||
|
||||
class TestAdverseSelection:
|
||||
def test_basic(self):
|
||||
mids_before = [100.0, 100.0, 100.0]
|
||||
mids_after = [100.5, 99.5, 100.0]
|
||||
sides = ["buy", "sell", "buy"]
|
||||
result = adverse_selection_ratio(mids_after, mids_before, sides)
|
||||
assert result["buy_adverse_bps"] != 0 or result["sell_adverse_bps"] != 0
|
||||
|
||||
|
||||
class TestLiquidationClustering:
|
||||
def test_no_clusters(self):
|
||||
times = list(range(0, 60000, 5000)) # spaced 5 sec apart
|
||||
result = liquidation_clustering(times)
|
||||
assert result["n_clusters"] == 0
|
||||
|
||||
def test_tight_clusters(self):
|
||||
times = [0, 100, 200, 300, 400, 50000, 50100, 50200, 50300]
|
||||
result = liquidation_clustering(times, window_sec=300)
|
||||
assert result["n_clusters"] >= 0 # may or may not cluster depending on mean interval
|
||||
|
||||
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
# Funding tests
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
|
||||
class TestFundingRegime:
|
||||
def test_neutral(self):
|
||||
rates = [0.000001] * 2000 # ~0.1% annual
|
||||
result = funding_regime(rates, window_hours=24, n_samples_per_hour=60)
|
||||
assert result["regime"] == "neutral"
|
||||
|
||||
def test_high_positive(self):
|
||||
rates = [0.0001] * 2000 # ~11% annual
|
||||
result = funding_regime(rates, window_hours=24, n_samples_per_hour=60)
|
||||
assert result["regime"] in ("positive", "high_positive")
|
||||
|
||||
|
||||
class TestFundingPredictability:
|
||||
def test_random(self):
|
||||
np.random.seed(42)
|
||||
rates = list(np.random.normal(0, 0.001, 100))
|
||||
result = funding_predictability(rates)
|
||||
assert len(result["autocorr"]) == 3
|
||||
|
||||
|
||||
class TestBasisSpread:
|
||||
def test_premium(self):
|
||||
perp = [105.0, 106.0, 107.0]
|
||||
spot = [100.0, 101.0, 102.0]
|
||||
result = basis_spread(perp, spot)
|
||||
assert result["current_basis_bps"] > 0
|
||||
assert result["mean_basis_bps"] > 0
|
||||
|
||||
|
||||
class TestBasisConvergence:
|
||||
def test_mean_reverting(self):
|
||||
basis = [50.0, 45.0, 40.0, 35.0, 30.0, 25.0] * 100 # decaying
|
||||
result = basis_convergence_speed(basis)
|
||||
assert result["ar1_coef"] > 0
|
||||
|
||||
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
# Signals tests
|
||||
# ═══════════════════════════════════════════════════════════
|
||||
|
||||
class TestCompositeSignal:
|
||||
def test_strong_buy(self):
|
||||
result = composite_signal(obi=0.5, trade_imbalance=0.3, vpin=0.1, funding_regime="negative")
|
||||
assert result["signal"] == "buy"
|
||||
assert result["confidence"] > 0
|
||||
|
||||
def test_strong_sell(self):
|
||||
result = composite_signal(obi=-0.5, trade_imbalance=-0.3, vpin=0.1, funding_regime="high_positive")
|
||||
assert result["signal"] == "sell"
|
||||
assert result["confidence"] > 0
|
||||
|
||||
def test_neutral(self):
|
||||
result = composite_signal(obi=0.05, trade_imbalance=0.0, vpin=0.0)
|
||||
assert result["signal"] == "neutral"
|
||||
assert result["confidence"] >= 0
|
||||
|
||||
def test_toxic_reduces_confidence(self):
|
||||
norm = composite_signal(obi=0.4, trade_imbalance=0.3, vpin=0.1)
|
||||
toxic = composite_signal(obi=0.4, trade_imbalance=0.3, vpin=0.8)
|
||||
assert toxic["confidence"] < norm["confidence"]
|
||||
|
||||
def test_wide_spread_reduces_confidence(self):
|
||||
tight = composite_signal(obi=0.4, trade_imbalance=0.3, spread_bps=1.0)
|
||||
wide = composite_signal(obi=0.4, trade_imbalance=0.3, spread_bps=30.0)
|
||||
assert wide["confidence"] < tight["confidence"]
|
||||
|
||||
|
||||
class TestDetectHFTRegime:
|
||||
def test_toxic(self):
|
||||
assert detect_hft_regime(obi_std=0.1, spread_mean_bps=15.0, trade_rate_per_sec=1.0, vpin=0.5) == "toxic"
|
||||
|
||||
def test_quiet(self):
|
||||
assert detect_hft_regime(obi_std=0.1, spread_mean_bps=2.0, trade_rate_per_sec=0.05, vpin=0.1) == "quiet"
|
||||
|
||||
def test_trending(self):
|
||||
assert detect_hft_regime(obi_std=0.4, spread_mean_bps=3.0, trade_rate_per_sec=2.0, vpin=0.1) == "trending"
|
||||
|
||||
def test_ranging(self):
|
||||
assert detect_hft_regime(obi_std=0.1, spread_mean_bps=2.0, trade_rate_per_sec=2.0, vpin=0.1) == "ranging"
|
||||
|
||||
|
||||
class TestSignalPipeline:
|
||||
def test_emit_no_data(self):
|
||||
p = SignalPipeline()
|
||||
result = p.emit()
|
||||
assert "signal" in result
|
||||
assert result["signal"] == "neutral"
|
||||
|
||||
def test_update_and_emit(self):
|
||||
p = SignalPipeline(obi_window=10, trade_window=10, vpin_volume_size=1.0, vpin_buckets=5)
|
||||
bids = {100.0: 5.0}
|
||||
asks = {102.0: 1.0}
|
||||
p.update_book(bids, asks)
|
||||
for _ in range(3):
|
||||
p.update_trade(101.5, 0.1, 101.0)
|
||||
for _ in range(1):
|
||||
p.update_trade(100.5, 0.1, 101.0)
|
||||
result = p.emit()
|
||||
assert result["signal"] in ("buy", "sell", "neutral")
|
||||
@@ -0,0 +1,123 @@
|
||||
"""
|
||||
Tests for data/normalizer.py — timestamp normalization and sequence gaps.
|
||||
"""
|
||||
import time
|
||||
|
||||
from data.normalizer import (
|
||||
normalize_timestamp,
|
||||
SequenceTracker,
|
||||
detect_sequence_gap,
|
||||
)
|
||||
|
||||
|
||||
def test_normalize_ms_timestamp():
|
||||
"""Millisecond timestamps pass through unchanged."""
|
||||
ts = 1705312800000
|
||||
assert normalize_timestamp(ts) == 1705312800000
|
||||
|
||||
|
||||
def test_normalize_seconds_timestamp():
|
||||
"""Second timestamps get multiplied by 1000."""
|
||||
ts = 1705312800
|
||||
result = normalize_timestamp(ts)
|
||||
assert result == 1705312800000
|
||||
|
||||
|
||||
def test_normalize_float_timestamp():
|
||||
"""Float seconds get multiplied."""
|
||||
ts = 1705312800.5
|
||||
result = normalize_timestamp(ts)
|
||||
assert result == 1705312800500
|
||||
|
||||
|
||||
def test_normalize_large_float_is_ms():
|
||||
"""A float > 1e12 is already in ms."""
|
||||
ts = 1705312800000.123
|
||||
result = normalize_timestamp(ts)
|
||||
assert result == 1705312800000
|
||||
|
||||
|
||||
def test_normalize_iso_string():
|
||||
"""ISO 8601 Z string converted to ms."""
|
||||
result = normalize_timestamp("2024-01-15T12:00:00.000Z")
|
||||
assert result == 1705320000000
|
||||
|
||||
|
||||
def test_normalize_iso_no_z():
|
||||
"""ISO string without trailing Z."""
|
||||
result = normalize_timestamp("2024-01-15T12:00:00")
|
||||
expected = 1705320000000
|
||||
assert abs(result - expected) < 1000
|
||||
|
||||
|
||||
def test_normalize_none_returns_now():
|
||||
"""None returns current time (within 1s)."""
|
||||
now_ms = int(time.time() * 1000)
|
||||
result = normalize_timestamp(None)
|
||||
assert abs(result - now_ms) < 2000
|
||||
|
||||
|
||||
def test_normalize_datetime():
|
||||
"""datetime object converted to ms."""
|
||||
from datetime import datetime
|
||||
dt = datetime(2024, 1, 15, 12, 0, 0)
|
||||
result = normalize_timestamp(dt)
|
||||
assert result == 1705320000000
|
||||
|
||||
|
||||
class TestSequenceTracker:
|
||||
def test_initial_no_gap(self):
|
||||
st = SequenceTracker()
|
||||
assert st.check("l2", "BTC", 100) is None
|
||||
|
||||
def test_consecutive_no_gap(self):
|
||||
st = SequenceTracker()
|
||||
st.check("l2", "BTC", 100)
|
||||
assert st.check("l2", "BTC", 101) is None
|
||||
assert st.check("l2", "BTC", 102) is None
|
||||
|
||||
def test_gap_detected(self):
|
||||
st = SequenceTracker()
|
||||
st.check("l2", "BTC", 100)
|
||||
gap = st.check("l2", "BTC", 105)
|
||||
assert gap is not None
|
||||
assert gap["gap_size"] == 4
|
||||
assert gap["expected"] == 101
|
||||
|
||||
def test_multiple_channels_independent(self):
|
||||
st = SequenceTracker()
|
||||
st.check("l2", "BTC", 100)
|
||||
st.check("trades", "BTC", 50)
|
||||
assert st.check("l2", "BTC", 101) is None
|
||||
assert st.check("trades", "BTC", 51) is None
|
||||
|
||||
def test_reset_clears_state(self):
|
||||
st = SequenceTracker()
|
||||
st.check("l2", "BTC", 100)
|
||||
st.reset("l2", "BTC")
|
||||
assert st.check("l2", "BTC", 200) is None # fresh start
|
||||
|
||||
def test_gap_counts_accumulate(self):
|
||||
st = SequenceTracker()
|
||||
st.check("l2", "BTC", 100)
|
||||
st.check("l2", "BTC", 105)
|
||||
st.check("l2", "BTC", 110)
|
||||
gaps = st.gap_counts
|
||||
assert gaps["l2:BTC"] == 8 # 4 + 4
|
||||
|
||||
|
||||
def test_detect_sequence_gap_empty():
|
||||
assert detect_sequence_gap(1, None) == 0
|
||||
|
||||
|
||||
def test_detect_sequence_gap_ok():
|
||||
assert detect_sequence_gap(102, 101) == 0
|
||||
|
||||
|
||||
def test_detect_sequence_gap_found():
|
||||
gap = detect_sequence_gap(105, 101)
|
||||
assert gap == 3
|
||||
|
||||
|
||||
def test_detect_sequence_gap_negative():
|
||||
assert detect_sequence_gap(100, 101) == -1 # dupe or reset
|
||||
@@ -0,0 +1,110 @@
|
||||
"""
|
||||
Tests for data/store.py — Parquet-based raw message storage.
|
||||
"""
|
||||
import json
|
||||
import os
|
||||
import tempfile
|
||||
import time
|
||||
|
||||
from data.store import RawMessageStore, read_range
|
||||
|
||||
|
||||
def _make_parquet_dir():
|
||||
"""Create a temporary directory for parquet storage."""
|
||||
return tempfile.mkdtemp()
|
||||
|
||||
|
||||
def test_store_basic_write_read():
|
||||
"""Write messages, flush, read them back."""
|
||||
tmpdir = tempfile.mkdtemp()
|
||||
store = RawMessageStore(data_dir=tmpdir, flush_interval_sec=0.5)
|
||||
store.start()
|
||||
|
||||
for i in range(10):
|
||||
store.push(
|
||||
channel="l2book",
|
||||
coin="BTC",
|
||||
exchange_ts=1718000000000 + i * 1000,
|
||||
payload={"type": "snapshot", "levels": [[{"px": "50000", "sz": "1.0"}], [{"px": "50001", "sz": "0.5"}]]},
|
||||
)
|
||||
|
||||
time.sleep(1.5)
|
||||
store.stop()
|
||||
assert store.total_written == 10
|
||||
|
||||
from datetime import date
|
||||
today = date.today().isoformat()
|
||||
rows = read_range(tmpdir, "l2book", "BTC", today, today)
|
||||
assert len(rows) == 10
|
||||
assert rows[0]["payload"]["type"] == "snapshot"
|
||||
assert rows[0]["exchange_ts"] == 1718000000000
|
||||
assert rows[-1]["exchange_ts"] == 1718000009000
|
||||
|
||||
|
||||
def test_store_multiple_channels():
|
||||
"""Write to different channels and verify partitioning."""
|
||||
tmpdir = tempfile.mkdtemp()
|
||||
store = RawMessageStore(data_dir=tmpdir, flush_interval_sec=0.5)
|
||||
store.start()
|
||||
|
||||
channels = ["l2book", "trades", "funding", "mark", "open_interest"]
|
||||
for ch in channels:
|
||||
for i in range(3):
|
||||
store.push(channel=ch, coin="BTC", exchange_ts=1718000000000 + i * 1000, payload={"ch": ch, "i": i})
|
||||
|
||||
time.sleep(1.5)
|
||||
store.stop()
|
||||
assert store.total_written == 15
|
||||
|
||||
from datetime import date
|
||||
today = date.today().isoformat()
|
||||
for ch in channels:
|
||||
rows = read_range(tmpdir, ch, "BTC", today, today)
|
||||
assert len(rows) == 3, f"Expected 3 rows for {ch}, got {len(rows)}"
|
||||
|
||||
|
||||
def test_store_append_to_existing():
|
||||
"""Write in two batches to same file — should append."""
|
||||
tmpdir = tempfile.mkdtemp()
|
||||
store = RawMessageStore(data_dir=tmpdir, flush_interval_sec=0.3)
|
||||
store.start()
|
||||
|
||||
for i in range(5):
|
||||
store.push(channel="trades", coin="ETH", exchange_ts=1718000000000 + i * 1000, payload={"batch": 1, "i": i})
|
||||
|
||||
time.sleep(1)
|
||||
store.stop()
|
||||
|
||||
store2 = RawMessageStore(data_dir=tmpdir, flush_interval_sec=0.3)
|
||||
store2.start()
|
||||
for i in range(5):
|
||||
store2.push(channel="trades", coin="ETH", exchange_ts=1718000005000 + i * 1000, payload={"batch": 2, "i": i})
|
||||
time.sleep(1)
|
||||
store2.stop()
|
||||
|
||||
from datetime import date
|
||||
today = date.today().isoformat()
|
||||
rows = read_range(tmpdir, "trades", "ETH", today, today)
|
||||
assert len(rows) == 10
|
||||
batches = [r["payload"]["batch"] for r in rows]
|
||||
assert batches == [1] * 5 + [2] * 5
|
||||
|
||||
|
||||
def test_store_queue_full_does_not_crash():
|
||||
"""Small queue — push many, ensure no crash."""
|
||||
tmpdir = tempfile.mkdtemp()
|
||||
store = RawMessageStore(data_dir=tmpdir, flush_interval_sec=60.0, max_queue_size=10)
|
||||
store.start()
|
||||
|
||||
for i in range(1000):
|
||||
store.push(channel="trades", coin="BTC", exchange_ts=1718000000000 + i, payload={"i": i})
|
||||
|
||||
store.stop()
|
||||
assert store.total_written >= 0 # some may be lost, but no crash
|
||||
|
||||
|
||||
def test_read_range_empty():
|
||||
"""Read a date range with no data returns empty list."""
|
||||
tmpdir = tempfile.mkdtemp()
|
||||
rows = read_range(tmpdir, "nonexistent", "BTC", "2020-01-01", "2020-01-02")
|
||||
assert rows == []
|
||||
Reference in New Issue
Block a user