273 lines
12 KiB
Python
273 lines
12 KiB
Python
|
|
"""
|
||
|
|
MALKHUT Asset Store — DuckDB file-backed storage.
|
||
|
|
|
||
|
|
High-performance columnar store for the system-wide asset universe.
|
||
|
|
Replaces in-memory dicts with DuckDB for persistence, query, and analytics.
|
||
|
|
|
||
|
|
Usage:
|
||
|
|
from malkhut.storage.asset_store import AssetStore
|
||
|
|
|
||
|
|
store = AssetStore() # opens malkhut_assets.duckdb
|
||
|
|
store.sync_from_profiles() # populate from in-memory AssetProfile dict
|
||
|
|
assets = store.query_assets(sector='layer1')
|
||
|
|
btc = store.get_asset('BTCUSDT')
|
||
|
|
"""
|
||
|
|
from __future__ import annotations
|
||
|
|
|
||
|
|
import os
|
||
|
|
from pathlib import Path
|
||
|
|
from typing import Any, Dict, List, Optional, Sequence, Tuple
|
||
|
|
|
||
|
|
import duckdb
|
||
|
|
|
||
|
|
|
||
|
|
_DEFAULT_DB_PATH = str(Path(__file__).resolve().parent / "malkhut_assets.duckdb")
|
||
|
|
|
||
|
|
|
||
|
|
class AssetStore:
|
||
|
|
"""DuckDB-backed asset universe store. Thread-safe reads, single-writer writes."""
|
||
|
|
|
||
|
|
def __init__(self, db_path: Optional[str] = None) -> None:
|
||
|
|
self.db_path = db_path or os.environ.get("MALKHUT_DUCKDB_PATH", _DEFAULT_DB_PATH)
|
||
|
|
self.conn = duckdb.connect(self.db_path)
|
||
|
|
self._ensure_schema()
|
||
|
|
|
||
|
|
def _ensure_schema(self) -> None:
|
||
|
|
self.conn.execute('''
|
||
|
|
CREATE TABLE IF NOT EXISTS exchanges (
|
||
|
|
exchange_id VARCHAR PRIMARY KEY,
|
||
|
|
display_name VARCHAR NOT NULL,
|
||
|
|
has_spot BOOLEAN DEFAULT true,
|
||
|
|
has_perps BOOLEAN DEFAULT true,
|
||
|
|
has_options BOOLEAN DEFAULT false,
|
||
|
|
api_base_url VARCHAR DEFAULT '',
|
||
|
|
ws_base_url VARCHAR DEFAULT '',
|
||
|
|
default_taker_fee_bps DOUBLE DEFAULT 0.0,
|
||
|
|
default_maker_fee_bps DOUBLE DEFAULT 0.0,
|
||
|
|
typical_latency_ms DOUBLE DEFAULT 0.0
|
||
|
|
)
|
||
|
|
''')
|
||
|
|
self.conn.execute('''
|
||
|
|
CREATE TABLE IF NOT EXISTS assets (
|
||
|
|
symbol VARCHAR PRIMARY KEY,
|
||
|
|
base_asset VARCHAR NOT NULL,
|
||
|
|
name VARCHAR NOT NULL,
|
||
|
|
unified_symbol VARCHAR NOT NULL,
|
||
|
|
quote_currency VARCHAR NOT NULL DEFAULT 'USDT',
|
||
|
|
coingecko_id VARCHAR DEFAULT '',
|
||
|
|
cmc_id INTEGER DEFAULT 0,
|
||
|
|
blockchain VARCHAR DEFAULT '',
|
||
|
|
contract_address VARCHAR DEFAULT '',
|
||
|
|
sectors VARCHAR[] NOT NULL,
|
||
|
|
token_roles VARCHAR[] NOT NULL,
|
||
|
|
supply_model VARCHAR NOT NULL,
|
||
|
|
consensus VARCHAR NOT NULL,
|
||
|
|
smart_contracts VARCHAR NOT NULL,
|
||
|
|
market_cap_tier VARCHAR NOT NULL,
|
||
|
|
volatility_profile VARCHAR NOT NULL,
|
||
|
|
liquidity_profile VARCHAR NOT NULL,
|
||
|
|
derivative_access VARCHAR NOT NULL,
|
||
|
|
tick_size DOUBLE NOT NULL,
|
||
|
|
lot_size DOUBLE NOT NULL,
|
||
|
|
price_decimals INTEGER NOT NULL,
|
||
|
|
maker_fee_bps DOUBLE NOT NULL,
|
||
|
|
taker_fee_bps DOUBLE NOT NULL,
|
||
|
|
typical_spread_bps DOUBLE NOT NULL,
|
||
|
|
typical_depth_usd DOUBLE NOT NULL,
|
||
|
|
typical_daily_volume_usd DOUBLE NOT NULL,
|
||
|
|
has_funding BOOLEAN DEFAULT false,
|
||
|
|
has_options BOOLEAN DEFAULT false
|
||
|
|
)
|
||
|
|
''')
|
||
|
|
self.conn.execute('''
|
||
|
|
CREATE TABLE IF NOT EXISTS asset_exchanges (
|
||
|
|
symbol VARCHAR NOT NULL,
|
||
|
|
exchange_id VARCHAR NOT NULL,
|
||
|
|
PRIMARY KEY (symbol, exchange_id)
|
||
|
|
)
|
||
|
|
''')
|
||
|
|
self.conn.execute('''
|
||
|
|
CREATE TABLE IF NOT EXISTS behavior_profiles (
|
||
|
|
symbol VARCHAR PRIMARY KEY,
|
||
|
|
template_name VARCHAR DEFAULT '',
|
||
|
|
reference_price DOUBLE DEFAULT 0.0,
|
||
|
|
depth_amplitude_usd DOUBLE NOT NULL,
|
||
|
|
depth_alpha DOUBLE NOT NULL,
|
||
|
|
depth_fragility DOUBLE NOT NULL,
|
||
|
|
depth_at_10bps_usd DOUBLE NOT NULL,
|
||
|
|
depth_at_100bps_usd DOUBLE NOT NULL,
|
||
|
|
spread_normal_bps DOUBLE NOT NULL,
|
||
|
|
spread_stress_mult DOUBLE NOT NULL,
|
||
|
|
flow_orders_per_sec DOUBLE NOT NULL,
|
||
|
|
flow_cancel_fill_ratio DOUBLE NOT NULL,
|
||
|
|
flow_median_order_usd DOUBLE NOT NULL,
|
||
|
|
flow_p99_order_usd DOUBLE NOT NULL,
|
||
|
|
vol_annualized_normal DOUBLE NOT NULL,
|
||
|
|
vol_annualized_crisis DOUBLE NOT NULL,
|
||
|
|
vol_garch_alpha DOUBLE NOT NULL,
|
||
|
|
vol_garch_beta DOUBLE NOT NULL,
|
||
|
|
vol_half_life_hours DOUBLE NOT NULL,
|
||
|
|
retail_ratio DOUBLE NOT NULL,
|
||
|
|
retail_inst_gap DOUBLE NOT NULL,
|
||
|
|
liq_oi_mcap_ratio DOUBLE NOT NULL,
|
||
|
|
liq_trigger_pct DOUBLE NOT NULL,
|
||
|
|
liq_speed VARCHAR NOT NULL,
|
||
|
|
liq_recovery VARCHAR NOT NULL,
|
||
|
|
bingx_spread_mult DOUBLE NOT NULL,
|
||
|
|
bingx_depth_ratio DOUBLE NOT NULL,
|
||
|
|
bingx_latency_ms DOUBLE NOT NULL
|
||
|
|
)
|
||
|
|
''')
|
||
|
|
self.conn.execute('CREATE INDEX IF NOT EXISTS idx_assets_blockchain ON assets(blockchain)')
|
||
|
|
self.conn.execute('CREATE INDEX IF NOT EXISTS idx_assets_coingecko ON assets(coingecko_id)')
|
||
|
|
self.conn.execute('CREATE INDEX IF NOT EXISTS idx_assets_cmc ON assets(cmc_id)')
|
||
|
|
self.conn.execute('CREATE INDEX IF NOT EXISTS idx_asset_exchanges_ex ON asset_exchanges(exchange_id)')
|
||
|
|
|
||
|
|
# ── Write operations ────────────────────────────────────────────
|
||
|
|
|
||
|
|
def upsert_exchange(self, ex: Any) -> None:
|
||
|
|
self.conn.execute('''
|
||
|
|
INSERT OR REPLACE INTO exchanges VALUES (?,?,?,?,?,?,?,?,?,?)
|
||
|
|
''', [
|
||
|
|
ex.exchange_id, ex.display_name, ex.has_spot, ex.has_perps,
|
||
|
|
ex.has_options, ex.api_base_url, ex.ws_base_url,
|
||
|
|
ex.default_taker_fee_bps, ex.default_maker_fee_bps, ex.typical_latency_ms,
|
||
|
|
])
|
||
|
|
|
||
|
|
def upsert_asset(self, p: Any) -> None:
|
||
|
|
self.conn.execute('''
|
||
|
|
INSERT OR REPLACE INTO assets VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
|
||
|
|
''', [
|
||
|
|
p.symbol, p.base_asset, p.name, p.unified_symbol, p.quote_currency,
|
||
|
|
p.coingecko_id, p.cmc_id, p.blockchain, p.contract_address,
|
||
|
|
list(p.sectors), list(p.token_roles),
|
||
|
|
p.supply_model.value if hasattr(p.supply_model, 'value') else str(p.supply_model),
|
||
|
|
p.consensus.value if hasattr(p.consensus, 'value') else str(p.consensus),
|
||
|
|
p.smart_contracts.value if hasattr(p.smart_contracts, 'value') else str(p.smart_contracts),
|
||
|
|
p.market_cap_tier.value if hasattr(p.market_cap_tier, 'value') else str(p.market_cap_tier),
|
||
|
|
p.volatility_profile.value if hasattr(p.volatility_profile, 'value') else str(p.volatility_profile),
|
||
|
|
p.liquidity_profile.value if hasattr(p.liquidity_profile, 'value') else str(p.liquidity_profile),
|
||
|
|
p.derivative_access.value if hasattr(p.derivative_access, 'value') else str(p.derivative_access),
|
||
|
|
p.tick_size, p.lot_size, p.price_decimals,
|
||
|
|
p.maker_fee_bps, p.taker_fee_bps,
|
||
|
|
p.typical_spread_bps, p.typical_depth_usd, p.typical_daily_volume_usd,
|
||
|
|
p.has_funding, p.has_options,
|
||
|
|
])
|
||
|
|
|
||
|
|
def upsert_asset_exchanges(self, symbol: str, exchanges: Tuple[str, ...]) -> None:
|
||
|
|
self.conn.execute('DELETE FROM asset_exchanges WHERE symbol = ?', [symbol])
|
||
|
|
for ex in exchanges:
|
||
|
|
self.conn.execute(
|
||
|
|
'INSERT INTO asset_exchanges (symbol, exchange_id) VALUES (?, ?)',
|
||
|
|
[symbol, ex],
|
||
|
|
)
|
||
|
|
|
||
|
|
def upsert_behavior(self, b: Any) -> None:
|
||
|
|
self.conn.execute('''
|
||
|
|
INSERT OR REPLACE INTO behavior_profiles VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
|
||
|
|
''', [
|
||
|
|
b.symbol, b.template_name, b.reference_price,
|
||
|
|
b.depth.amplitude_usd, b.depth.alpha, b.depth.fragility_factor,
|
||
|
|
b.depth.depth_at_10bps_usd, b.depth.depth_at_100bps_usd,
|
||
|
|
b.spread.normal_bps, b.spread.stress_multiplier,
|
||
|
|
b.flow.orders_per_sec_normal, b.flow.cancel_fill_ratio,
|
||
|
|
b.flow.median_order_usd, b.flow.p99_order_usd,
|
||
|
|
b.vol.annualized_normal, b.vol.annualized_crisis,
|
||
|
|
b.vol.garch_alpha, b.vol.garch_beta, b.vol.half_life_hours,
|
||
|
|
b.retail.ratio, b.retail.inst_gap,
|
||
|
|
b.liquidation.oi_mcap_ratio, b.liquidation.trigger_pct,
|
||
|
|
b.liquidation.speed, b.liquidation.recovery,
|
||
|
|
b.bingx.spread_mult, b.bingx.depth_ratio, b.bingx.latency_ms,
|
||
|
|
])
|
||
|
|
|
||
|
|
# ── Read operations ─────────────────────────────────────────────
|
||
|
|
|
||
|
|
def get_asset(self, symbol: str) -> Optional[dict]:
|
||
|
|
row = self.conn.execute(
|
||
|
|
'SELECT * FROM assets WHERE symbol = ?', [symbol]
|
||
|
|
).fetchone()
|
||
|
|
if row is None:
|
||
|
|
return None
|
||
|
|
cols = [d[0] for d in self.conn.description]
|
||
|
|
return dict(zip(cols, row))
|
||
|
|
|
||
|
|
def get_asset_exchanges(self, symbol: str) -> List[str]:
|
||
|
|
rows = self.conn.execute(
|
||
|
|
'SELECT exchange_id FROM asset_exchanges WHERE symbol = ? ORDER BY exchange_id',
|
||
|
|
[symbol],
|
||
|
|
).fetchall()
|
||
|
|
return [r[0] for r in rows]
|
||
|
|
|
||
|
|
def query_assets(self, **filters: Any) -> List[dict]:
|
||
|
|
"""Query assets with optional WHERE filters. Arrays use list_contains."""
|
||
|
|
where_parts = []
|
||
|
|
params = []
|
||
|
|
for key, val in filters.items():
|
||
|
|
if isinstance(val, list):
|
||
|
|
# DuckDB: check if array column contains any of the values
|
||
|
|
placeholders = ", ".join(["?" for _ in val])
|
||
|
|
where_parts.append(f"list_has_any({key}, ARRAY[{placeholders}])")
|
||
|
|
params.extend(val)
|
||
|
|
elif isinstance(val, str):
|
||
|
|
where_parts.append(f"{key} = ?")
|
||
|
|
params.append(val)
|
||
|
|
elif isinstance(val, (int, float)):
|
||
|
|
where_parts.append(f"{key} = ?")
|
||
|
|
params.append(val)
|
||
|
|
elif isinstance(val, bool):
|
||
|
|
where_parts.append(f"{key} = ?")
|
||
|
|
params.append(val)
|
||
|
|
where_clause = " AND ".join(where_parts) if where_parts else "1=1"
|
||
|
|
rows = self.conn.execute(
|
||
|
|
f'SELECT * FROM assets WHERE {where_clause} ORDER BY symbol', params
|
||
|
|
).fetchall()
|
||
|
|
cols = [d[0] for d in self.conn.description]
|
||
|
|
return [dict(zip(cols, row)) for row in rows]
|
||
|
|
|
||
|
|
def symbols_for_exchange(self, exchange_id: str) -> List[str]:
|
||
|
|
rows = self.conn.execute('''
|
||
|
|
SELECT a.symbol FROM assets a
|
||
|
|
JOIN asset_exchanges ae ON a.symbol = ae.symbol
|
||
|
|
WHERE LOWER(ae.exchange_id) = LOWER(?)
|
||
|
|
ORDER BY a.symbol
|
||
|
|
''', [exchange_id]).fetchall()
|
||
|
|
return [r[0] for r in rows]
|
||
|
|
|
||
|
|
def assets_on_blockchain(self, blockchain: str) -> List[str]:
|
||
|
|
rows = self.conn.execute(
|
||
|
|
'SELECT symbol FROM assets WHERE blockchain = ? ORDER BY symbol',
|
||
|
|
[blockchain],
|
||
|
|
).fetchall()
|
||
|
|
return [r[0] for r in rows]
|
||
|
|
|
||
|
|
def asset_count(self) -> int:
|
||
|
|
return self.conn.execute('SELECT COUNT(*) FROM assets').fetchone()[0]
|
||
|
|
|
||
|
|
def exchange_count(self) -> int:
|
||
|
|
return self.conn.execute('SELECT COUNT(*) FROM exchanges').fetchone()[0]
|
||
|
|
|
||
|
|
# ── Sync from Python dicts ──────────────────────────────────────
|
||
|
|
|
||
|
|
def sync_from_profiles(self) -> int:
|
||
|
|
"""Populate DuckDB from in-memory ASSET_PROFILES + ASSET_BEHAVIORS + EXCHANGE_PROFILES."""
|
||
|
|
from malkhut.training.asset_classification import (
|
||
|
|
ASSET_PROFILES, EXCHANGE_PROFILES,
|
||
|
|
)
|
||
|
|
from malkhut.training.asset_behavior import ASSET_BEHAVIORS
|
||
|
|
|
||
|
|
count = 0
|
||
|
|
for ex in EXCHANGE_PROFILES.values():
|
||
|
|
self.upsert_exchange(ex)
|
||
|
|
for p in ASSET_PROFILES.values():
|
||
|
|
self.upsert_asset(p)
|
||
|
|
self.upsert_asset_exchanges(p.symbol, p.exchanges)
|
||
|
|
count += 1
|
||
|
|
for b in ASSET_BEHAVIORS.values():
|
||
|
|
if b.symbol in ASSET_PROFILES:
|
||
|
|
self.upsert_behavior(b)
|
||
|
|
self.conn.commit()
|
||
|
|
return count
|
||
|
|
|
||
|
|
def close(self) -> None:
|
||
|
|
self.conn.close()
|