fcfc136384
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)
99 lines
3.6 KiB
Python
99 lines
3.6 KiB
Python
"""
|
|
Exchange vs signal latency tracking.
|
|
|
|
Measures:
|
|
1. Exchange transport latency: exchange_ts → local receipt time
|
|
2. Signal computation latency: data receipt → signal generated
|
|
3. Order latency: signal → order accepted on exchange
|
|
4. Round-trip latency: signal → fill confirmation
|
|
|
|
Each metric is tracked as a rolling window with percentiles.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import time
|
|
from collections import deque
|
|
from typing import Optional
|
|
|
|
|
|
class LatencyTracker:
|
|
"""Track exchange and signal latencies with rolling percentiles."""
|
|
|
|
def __init__(self, window_seconds: float = 300.0, max_samples: int = 10000):
|
|
self._window = window_seconds
|
|
self._transport: deque[tuple[float, float]] = deque(maxlen=max_samples) # (time, ms)
|
|
self._signal: deque[tuple[float, float]] = deque(maxlen=max_samples)
|
|
self._order: deque[tuple[float, float]] = deque(maxlen=max_samples)
|
|
self._roundtrip: deque[tuple[float, float]] = deque(maxlen=max_samples)
|
|
|
|
def record_transport(self, exchange_ts_ms: int, local_ts: float | None = None):
|
|
"""Exchange timestamp → local receipt (ms)."""
|
|
local = local_ts or time.time()
|
|
lat = (local * 1000) - exchange_ts_ms
|
|
if 0 <= lat < 300_000: # Ignore clock skew > 5 min
|
|
self._transport.append((time.time(), lat))
|
|
|
|
def record_signal(self, duration_ms: float):
|
|
"""Time from data receipt to signal generation (ms)."""
|
|
if duration_ms >= 0:
|
|
self._signal.append((time.time(), duration_ms))
|
|
|
|
def record_order(self, duration_ms: float):
|
|
"""Signal generation → order accepted on exchange (ms)."""
|
|
if duration_ms >= 0:
|
|
self._order.append((time.time(), duration_ms))
|
|
|
|
def record_roundtrip(self, duration_ms: float):
|
|
"""Signal generation → fill confirmed (ms)."""
|
|
if duration_ms >= 0:
|
|
self._roundtrip.append((time.time(), duration_ms))
|
|
|
|
# ── Stats ──────────────────────────────────────────────────
|
|
|
|
def stats(self) -> dict:
|
|
return {
|
|
"transport_ms": self._percentiles(self._transport),
|
|
"signal_ms": self._percentiles(self._signal),
|
|
"order_ms": self._percentiles(self._order),
|
|
"roundtrip_ms": self._percentiles(self._roundtrip),
|
|
}
|
|
|
|
def summary(self) -> dict:
|
|
"""Compact summary: just p50/p99 for each metric."""
|
|
s = self.stats()
|
|
out = {}
|
|
for key, pct in s.items():
|
|
out[key] = {"p50": pct.get("p50", 0), "p99": pct.get("p99", 0)}
|
|
return out
|
|
|
|
# ── Internals ──────────────────────────────────────────────
|
|
|
|
def _prune(self, buffer: deque):
|
|
cutoff = time.time() - self._window
|
|
while buffer and buffer[0][0] < cutoff:
|
|
buffer.popleft()
|
|
|
|
def _percentiles(self, buffer: deque) -> dict:
|
|
self._prune(buffer)
|
|
if not buffer:
|
|
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(_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,
|
|
}
|