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)
127 lines
4.0 KiB
Python
127 lines
4.0 KiB
Python
"""
|
|
Timestamp normalization and sequence gap detection for market data.
|
|
|
|
Exchange timestamps come in various formats (ms since epoch, ISO strings,
|
|
exchange-specific formats). This module normalizes them to a consistent
|
|
int64 milliseconds-since-epoch.
|
|
|
|
Gap detection tracks per-channel per-coin sequence numbers and flags
|
|
missing messages so order books can be re-snapshotted.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import time
|
|
from datetime import datetime, timezone
|
|
|
|
# ── Timestamp normalization ───────────────────────────────────
|
|
|
|
def normalize_timestamp(ts, source: str = "hl") -> int:
|
|
"""Normalize a timestamp to int64 milliseconds since epoch.
|
|
|
|
Args:
|
|
ts: raw timestamp — can be int (ms), float (seconds), str (ISO 8601)
|
|
source: 'hl' (Hyperliquid), 'binance', 'bybit', 'okx', 'coinbase', 'deribit'
|
|
|
|
Returns int64 milliseconds since epoch.
|
|
"""
|
|
if ts is None:
|
|
return int(time.time() * 1000)
|
|
|
|
if isinstance(ts, (int, float)):
|
|
if ts > 1_000_000_000_000:
|
|
return int(ts) # already ms
|
|
if ts > 1_000_000_000:
|
|
return int(ts * 1000) # seconds → ms
|
|
return int(ts * 1000) # fractional seconds → ms
|
|
|
|
if isinstance(ts, str):
|
|
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)
|
|
|
|
|
|
def _parse_iso_ms(s: str) -> int:
|
|
for fmt in [
|
|
"%Y-%m-%dT%H:%M:%S.%fZ",
|
|
"%Y-%m-%dT%H:%M:%S.%f",
|
|
"%Y-%m-%dT%H:%M:%SZ",
|
|
"%Y-%m-%dT%H:%M:%S",
|
|
"%Y-%m-%d %H:%M:%S.%f",
|
|
"%Y-%m-%d %H:%M:%S",
|
|
]:
|
|
try:
|
|
dt = datetime.strptime(s.replace("+00:00", "").rstrip("Z"), fmt)
|
|
if dt.tzinfo is None:
|
|
dt = dt.replace(tzinfo=timezone.utc)
|
|
return int(dt.timestamp() * 1000)
|
|
except ValueError:
|
|
continue
|
|
return int(time.time() * 1000)
|
|
|
|
|
|
# ── Sequence gap detection ────────────────────────────────────
|
|
|
|
class SequenceTracker:
|
|
"""Track per-channel per-coin sequence numbers and detect gaps.
|
|
|
|
Usage:
|
|
tracker = SequenceTracker()
|
|
gap = tracker.check("l2book", "BTC", seq_num=1042)
|
|
if gap:
|
|
print(f"Gap detected: expected {gap['expected']}, got {gap['got']}")
|
|
"""
|
|
|
|
def __init__(self):
|
|
self._state: dict[str, int] = {} # key = "channel:coin", value = last_seq
|
|
self._gap_count: dict[str, int] = {}
|
|
|
|
def check(self, channel: str, coin: str, seq_num: int) -> dict | None:
|
|
"""Check for sequence gap. Returns None if ok, dict if gap."""
|
|
key = f"{channel}:{coin}"
|
|
last = self._state.get(key)
|
|
|
|
if last is None:
|
|
self._state[key] = seq_num
|
|
return None
|
|
|
|
expected = last + 1
|
|
if seq_num == expected or seq_num > expected:
|
|
self._state[key] = seq_num
|
|
if seq_num > expected:
|
|
gap_size = seq_num - expected
|
|
self._gap_count[key] = self._gap_count.get(key, 0) + gap_size
|
|
return {"key": key, "expected": expected, "got": seq_num, "gap_size": gap_size}
|
|
return None
|
|
|
|
return None
|
|
|
|
def reset(self, channel: str, coin: str):
|
|
"""Reset tracker (call after re-snapshot)."""
|
|
self._state.pop(f"{channel}:{coin}", None)
|
|
|
|
@property
|
|
def gap_counts(self) -> dict[str, int]:
|
|
return dict(self._gap_count)
|
|
|
|
|
|
def detect_sequence_gap(
|
|
current_seq: int,
|
|
last_seq: int | None,
|
|
max_gap: int = 10,
|
|
) -> int:
|
|
"""Return gap size. 0 = ok, >0 = gap count, -1 = negative gap (dupe/reset)."""
|
|
if last_seq is None:
|
|
return 0
|
|
diff = current_seq - last_seq
|
|
if diff == 1:
|
|
return 0
|
|
if diff > 1:
|
|
return min(diff - 1, max_gap * 100)
|
|
return -1 # duplicate or reset
|