0446443d36
New strategies: - Cross-Sectional Momentum: long top-N, short bottom-N across HL universe - Spot-Perp Basis Arbitrage: delta-neutral spot vs perp price gap trading - Regime-Switching Ensemble: dynamically allocates strategies by market regime - Portfolio Construction: risk parity, vol targeting, correlation penalty Infrastructure: - DuckDBDataProvider: real tick/candle data for backtests (replaces synthetic) - Walk-Forward Validation: systematic IS/OOS across all 12 strategies - 3 Jupyter research notebooks (EDA, strategy research, portfolio) Pipeline integration: - deploy.py registry, sweep_runner, vbt_runner all updated - fee_tiers support for new strategies - All modules syntax-validated and import-tested
290 lines
9.0 KiB
Python
290 lines
9.0 KiB
Python
"""
|
|
DuckDB-powered data provider for backtesting with real Hyperliquid data.
|
|
|
|
Replaces the REST-based HyperliquidDataProvider with local DuckDB queries
|
|
for fast historical backtesting on authentic market data. Falls back to
|
|
REST API when DuckDB isn't available or data is stale.
|
|
|
|
Provides:
|
|
- OHLCV candle construction from tick trades
|
|
- Multi-asset parallel fetching
|
|
- Pre-computed rollups (microprice, OFI, VPIN) at 1s/1m resolution
|
|
- Trade-level data for Hurst/VPIN backtests
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Optional
|
|
|
|
import numpy as np
|
|
import pandas as pd
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class DuckDBProvider:
|
|
"""Fast local data provider backed by DuckDB tick database."""
|
|
|
|
def __init__(self, db_path: str = "data/normalized/ftdt_tick.db"):
|
|
self._db_path = Path(db_path)
|
|
self._conn = None
|
|
self._available = self._db_path.exists()
|
|
self._rest_provider = None
|
|
|
|
@property
|
|
def available(self) -> bool:
|
|
return self._available
|
|
|
|
@property
|
|
def conn(self):
|
|
if self._conn is None and self._available:
|
|
import duckdb
|
|
self._conn = duckdb.connect(str(self._db_path))
|
|
return self._conn
|
|
|
|
def _ensure_rest(self):
|
|
if self._rest_provider is None:
|
|
from framework.data import HyperliquidDataProvider
|
|
self._rest_provider = HyperliquidDataProvider(testnet=False)
|
|
return self._rest_provider
|
|
|
|
def fetch_candles(
|
|
self,
|
|
coin: str,
|
|
interval: str = "1h",
|
|
start_ms: Optional[int] = None,
|
|
end_ms: Optional[int] = None,
|
|
limit: int = 5000,
|
|
) -> pd.DataFrame:
|
|
"""Fetch OHLCV candles, preferring DuckDB over REST API."""
|
|
if self._available:
|
|
df = self._fetch_candles_duckdb(coin, interval, start_ms, end_ms, limit)
|
|
if not df.empty:
|
|
return df
|
|
return self._ensure_rest().fetch_candles(coin, interval, start_ms, end_ms, limit)
|
|
|
|
def _fetch_candles_duckdb(
|
|
self,
|
|
coin: str,
|
|
interval: str,
|
|
start_ms: Optional[int] = None,
|
|
end_ms: Optional[int] = None,
|
|
limit: int = 5000,
|
|
) -> pd.DataFrame:
|
|
"""Build OHLCV candles from DuckDB trades table."""
|
|
interval_ms = {
|
|
"1m": 60_000, "5m": 300_000, "15m": 900_000, "30m": 1_800_000,
|
|
"1h": 3_600_000, "4h": 14_400_000, "1d": 86_400_000,
|
|
}.get(interval, 3_600_000)
|
|
|
|
if end_ms is None:
|
|
import time
|
|
end_ms = int(time.time() * 1000)
|
|
if start_ms is None:
|
|
start_ms = end_ms - limit * interval_ms
|
|
|
|
query = """
|
|
SELECT
|
|
(exchange_ts_ms / $interval_ms)::BIGINT * $interval_ms AS bucket_ts,
|
|
MIN(price) AS low,
|
|
MAX(price) AS high,
|
|
FIRST(price) AS open,
|
|
LAST(price) AS close,
|
|
SUM(size * price) AS volume
|
|
FROM trades
|
|
WHERE coin = $coin
|
|
AND exchange_ts_ms >= $start
|
|
AND exchange_ts_ms < $end
|
|
GROUP BY bucket_ts
|
|
ORDER BY bucket_ts
|
|
LIMIT $limit
|
|
"""
|
|
try:
|
|
result = self.conn.execute(query, {
|
|
"interval_ms": interval_ms,
|
|
"coin": coin.upper(),
|
|
"start": start_ms,
|
|
"end": end_ms,
|
|
"limit": limit,
|
|
}).fetchdf()
|
|
|
|
if result.empty:
|
|
return pd.DataFrame(columns=["open", "high", "low", "close", "volume", "timestamp"])
|
|
|
|
result["timestamp"] = pd.to_datetime(result["bucket_ts"], unit="ms", utc=True)
|
|
result = result[["open", "high", "low", "close", "volume", "timestamp"]]
|
|
result.set_index("timestamp", inplace=True)
|
|
result.sort_index(inplace=True)
|
|
return result.astype({k: float for k in ["open", "high", "low", "close", "volume"]})
|
|
except Exception as e:
|
|
logger.debug("DuckDB candle query failed: %s", e)
|
|
return pd.DataFrame()
|
|
|
|
def fetch_multi_candles(
|
|
self,
|
|
coins: list[str],
|
|
interval: str = "1h",
|
|
limit: int = 5000,
|
|
) -> dict[str, pd.DataFrame]:
|
|
"""Fetch candles for multiple coins. Uses DuckDB when available."""
|
|
results = {}
|
|
for coin in coins:
|
|
df = self.fetch_candles(coin, interval=interval, limit=limit)
|
|
if not df.empty:
|
|
results[coin] = df
|
|
return results
|
|
|
|
def fetch_trades(
|
|
self,
|
|
coin: str,
|
|
start_ms: Optional[int] = None,
|
|
end_ms: Optional[int] = None,
|
|
limit: int = 100_000,
|
|
) -> pd.DataFrame:
|
|
"""Fetch individual trades for Hurst/VPIN backtests."""
|
|
if not self._available:
|
|
return pd.DataFrame()
|
|
|
|
if end_ms is None:
|
|
import time
|
|
end_ms = int(time.time() * 1000)
|
|
if start_ms is None:
|
|
start_ms = end_ms - 24 * 3600 * 1000 # Default: 1 day
|
|
|
|
query = """
|
|
SELECT exchange_ts_ms, price, size, aggressor
|
|
FROM trades
|
|
WHERE coin = $coin
|
|
AND exchange_ts_ms >= $start
|
|
AND exchange_ts_ms < $end
|
|
ORDER BY exchange_ts_ms
|
|
LIMIT $limit
|
|
"""
|
|
try:
|
|
result = self.conn.execute(query, {
|
|
"coin": coin.upper(),
|
|
"start": start_ms,
|
|
"end": end_ms,
|
|
"limit": limit,
|
|
}).fetchdf()
|
|
return result
|
|
except Exception:
|
|
return pd.DataFrame()
|
|
|
|
def fetch_micro_rollup(
|
|
self,
|
|
coin: str,
|
|
start_ms: Optional[int] = None,
|
|
end_ms: Optional[int] = None,
|
|
limit: int = 100_000,
|
|
) -> pd.DataFrame:
|
|
"""Fetch pre-computed 1s microstructural rollups (microprice, OFI, trade imbalance)."""
|
|
if not self._available:
|
|
return pd.DataFrame()
|
|
|
|
if end_ms is None:
|
|
import time
|
|
end_ms = int(time.time() * 1000)
|
|
if start_ms is None:
|
|
start_ms = end_ms - 24 * 3600 * 1000
|
|
|
|
query = """
|
|
SELECT ts_1s, coin, mid_price, microprice, spread_bps,
|
|
buy_volume, sell_volume, trade_count, trade_imbalance
|
|
FROM micro_rollup_1s
|
|
WHERE coin = $coin
|
|
AND ts_1s >= $start
|
|
AND ts_1s < $end
|
|
ORDER BY ts_1s
|
|
LIMIT $limit
|
|
"""
|
|
try:
|
|
return self.conn.execute(query, {
|
|
"coin": coin.upper(),
|
|
"start": start_ms,
|
|
"end": end_ms,
|
|
"limit": limit,
|
|
}).fetchdf()
|
|
except Exception:
|
|
return pd.DataFrame()
|
|
|
|
def get_available_coins(self) -> list[str]:
|
|
"""List coins with data in DuckDB."""
|
|
if not self._available:
|
|
return ["BTC", "ETH"]
|
|
try:
|
|
result = self.conn.execute(
|
|
"SELECT DISTINCT coin FROM l2_snapshots ORDER BY coin"
|
|
).fetchall()
|
|
return [r[0] for r in result]
|
|
except Exception:
|
|
return ["BTC", "ETH"]
|
|
|
|
def get_data_range(self) -> tuple[int, int]:
|
|
"""Get min/max timestamps in database."""
|
|
if not self._available:
|
|
import time
|
|
now = int(time.time() * 1000)
|
|
return now - 30 * 24 * 3600 * 1000, now
|
|
try:
|
|
result = self.conn.execute(
|
|
"SELECT MIN(exchange_ts_ms), MAX(exchange_ts_ms) FROM l2_snapshots"
|
|
).fetchone()
|
|
return result[0] or 0, result[1] or 0
|
|
except Exception:
|
|
return 0, 0
|
|
|
|
def fetch_funding(
|
|
self,
|
|
coin: str,
|
|
start_ms: Optional[int] = None,
|
|
end_ms: Optional[int] = None,
|
|
limit: int = 10000,
|
|
) -> pd.DataFrame:
|
|
"""Fetch funding rate history."""
|
|
if not self._available:
|
|
return pd.DataFrame()
|
|
|
|
if end_ms is None:
|
|
import time
|
|
end_ms = int(time.time() * 1000)
|
|
if start_ms is None:
|
|
start_ms = end_ms - 30 * 24 * 3600 * 1000
|
|
|
|
query = """
|
|
SELECT exchange_ts_ms, coin, funding_rate, mark_px, annual_apr
|
|
FROM funding
|
|
WHERE coin = $coin
|
|
AND exchange_ts_ms >= $start
|
|
AND exchange_ts_ms < $end
|
|
ORDER BY exchange_ts_ms
|
|
LIMIT $limit
|
|
"""
|
|
try:
|
|
return self.conn.execute(query, {
|
|
"coin": coin.upper(),
|
|
"start": start_ms,
|
|
"end": end_ms,
|
|
"limit": limit,
|
|
}).fetchdf()
|
|
except Exception:
|
|
return pd.DataFrame()
|
|
|
|
def stats(self) -> dict:
|
|
"""Database statistics."""
|
|
if not self._available:
|
|
return {"status": "unavailable"}
|
|
from data.duckdb_load import DuckDBLoader
|
|
loader = DuckDBLoader(str(self._db_path))
|
|
try:
|
|
return loader.stats()
|
|
finally:
|
|
loader.close()
|
|
|
|
def close(self):
|
|
if self._conn:
|
|
self._conn.close()
|
|
self._conn = None
|