Files
ramseshk e58c5951b7 feat: Phase 5 — integration layer (analytics pipeline, production node v2, CLI) + 12 tests
live/integrator.py (AnalyticsPipeline):
  Real-time pipeline: data → microstructure → signals.
  Accumulates book snapshots + trades, computes OBI, VPIN, microprice,
  spread, depth, trade imbalance, HFT regime, and emits composite
  signal with confidence and breakdown. Per-coin isolation.

live/node_v2.py (ProductionNode):
  Rebuilt production node integrating ALL Phase 1-4 modules:
  - REST data fetching (order book, mark prices, funding rates)
  - AnalyticsPipeline per coin for real-time microstructure signals
  - Treasury for position/capital/PnL/breaker management
  - ToxicityFilter integration via HlMakerPool makers
  - HlMakerPool for per-coin A-S quoting
  - CrossVenueMonitor, FundingBasisMonitor, LiquidationRiskOverlay
  - Paper trading with probabilistic fill simulation
  - Dashboard metrics JSON output (equity, treasury, analytics, maker)
  - Periodic status logging

cli.py (unified CLI):
  Subcommands integrating all modules:
    collect   — Run Hyperliquid data collector to Parquet
    analyze   — Run microstructure analytics on stored data
    simulate  — Run market-making simulator on stored data
    run       — Start production trading node (paper or live)
    backtest  — Run VectorBT backtest

12 integration tests (all pass):
  - AnalyticsPipeline: empty, book, trade, VPIN, emit, regime, isolation
  - ProductionNode: creation, tick cycle (3 ticks), metrics JSON output
  - CLI: import verification

Total test suite: 184 tests, all passing.
2026-08-07 14:54:47 +08:00

145 lines
4.4 KiB
Python

"""
Integration tests for the analytics pipeline and production node.
"""
from live.integrator import AnalyticsPipeline
class TestAnalyticsPipeline:
def test_empty_pipeline(self):
p = AnalyticsPipeline()
result = p.emit()
assert result["signal"] == "neutral"
assert result["mid"] == 0
def test_book_update_sets_mid(self):
p = AnalyticsPipeline()
p.update_book({50000.0: 1.0, 49999.0: 2.0}, {50002.0: 1.0, 50003.0: 0.5})
assert p.mid == 50001.0
assert p.obi != 0
assert p.spread_bps > 0
def test_trade_updates_vpin(self):
p = AnalyticsPipeline()
p.update_book({50000.0: 5.0, 49999.0: 5.0}, {50001.0: 5.0, 50002.0: 5.0})
for _ in range(100):
p.update_trade(50001.0, 0.01, 50000.5) # buys
for _ in range(50):
p.update_trade(50000.0, 0.005, 50000.5) # sells
assert p.vpin >= 0
assert p.trade_imbalance() > 0 # more buys
def test_emi_of(self):
p = AnalyticsPipeline()
p.update_book({50000.0: 1.0, 49999.0: 2.0}, {50002.0: 1.0, 50003.0: 0.5})
assert p.ema_obi(alpha=0.5) != 0
def test_full_emit(self):
p = AnalyticsPipeline()
bids = {100.0: 1.0, 99.0: 2.0}
asks = {102.0: 1.0, 103.0: 3.0}
p.update_book(bids, asks)
for _ in range(20):
p.update_trade(101.0, 0.1, 101.0)
result = p.emit(funding_regime="neutral")
assert "signal" in result
assert "confidence" in result
assert "vpin" in result
assert "obi" in result
assert "hft_regime" in result
assert "breakdown" in result
def test_hft_regime_detect(self):
p = AnalyticsPipeline()
bids = {100.0: 1.0}
asks = {102.0: 1.0}
p.update_book(bids, asks)
regime = p.hft_regime()
assert regime in ("trending", "ranging", "toxic", "quiet")
def test_pipeline_per_coin_isolation(self):
btc = AnalyticsPipeline()
eth = AnalyticsPipeline()
btc.update_book({50000.0: 1.0}, {50002.0: 1.0})
eth.update_book({3000.0: 1.0}, {3002.0: 1.0})
assert btc.mid > 40000
assert eth.mid < 10000
def test_trade_count_tracking(self):
p = AnalyticsPipeline()
p.update_book({100.0: 1.0}, {102.0: 1.0})
for _ in range(5):
p.update_trade(101.0, 0.1)
assert p.emit()["trade_count"] == 5
class TestNodeV2Smoke:
def test_node_creation(self):
from live.node_v2 import ProductionNode
node = ProductionNode(
coins=["BTC"],
testnet=True,
mode="paper",
max_position_per_coin=0.001,
base_quote_size=0.0001,
)
assert node is not None
def test_node_start_stop(self):
import asyncio
from live.node_v2 import ProductionNode
async def _test():
node = ProductionNode(
coins=["BTC"],
testnet=True,
mode="paper",
tick_interval_sec=0.1,
max_position_per_coin=0.001,
base_quote_size=0.0001,
)
await node.start()
for _ in range(3):
await node._tick_cycle()
await node.stop()
assert node._tick > 0
asyncio.run(_test())
def test_metrics_written(self):
import asyncio, json, tempfile, os, time
from live.node_v2 import ProductionNode
f = tempfile.NamedTemporaryFile(delete=False, suffix=".json")
f.close()
async def _test():
node = ProductionNode(
coins=["BTC"],
testnet=True,
mode="paper",
tick_interval_sec=0.1,
max_position_per_coin=0.001,
base_quote_size=0.0001,
metrics_file=f.name,
)
await node.start()
for _ in range(2):
await node._tick_cycle()
await node.stop()
assert os.path.exists(f.name)
data = json.load(open(f.name))
assert "treasury" in data
assert "analytics" in data
assert "equity_history" in data
assert "maker" in data
asyncio.run(_test())
os.unlink(f.name)
class TestCLISmoke:
def test_cli_import(self):
import cli
assert cli.main is not None