malkhut: DuckDB file-backed asset store — schema, sync, queries, persistence
DuckDB store (asset_store.py): - 4 tables: assets (28 cols), exchanges (10 cols), asset_exchanges (junction), behavior_profiles (28 cols) - Schema with indexes on blockchain, coingecko_id, cmc_id, exchange_id - sync_from_profiles(): populate from in-memory dicts in one call - Query API: query_assets(**filters), get_asset(), symbols_for_exchange(), assets_on_blockchain(), asset_count(), exchange_count() - Case-insensitive exchange lookup - File-backed persistence: data survives connection close/reopen - Performance: <1s sync, 100 queries in <1s, 1000 gets in <1s 30 tests covering: - Schema creation and table structure (3 tests) - Sync from profiles (6 tests, idempotent) - Asset CRUD (7 tests, query by blockchain/coingecko/sector) - Exchange mapping (5 tests, case-insensitive) - Behavior profiles (4 tests) - Persistence across connections (2 tests) - Performance baseline (3 tests) Total: 1276 tests across 50 files, all green, zero regressions.
This commit is contained in:
@@ -118,7 +118,8 @@ MALKHUT/
|
||||
│ │ └── control_plane.py # CONTROL_PLANE shared memory region
|
||||
│ ├── storage/
|
||||
│ │ ├── __init__.py
|
||||
│ │ └── ch_store.py # ClickHouse persistence (5 tables)
|
||||
│ │ ├── ch_store.py # ClickHouse persistence (5 tables)
|
||||
│ │ └── asset_store.py # DuckDB file-backed asset universe store
|
||||
│ ├── execution/
|
||||
│ │ ├── __init__.py
|
||||
│ │ └── asex_integration.py # ASEx validate-before-mutate kernel
|
||||
@@ -398,7 +399,7 @@ simple doctrinal tick-exits (C11) ship first via T19 step 3; MALKHUT supersedes
|
||||
|
||||
## DEVELOPMENT STATUS (2026-07-11)
|
||||
|
||||
**1189 test functions. 47 test files. All green. 0 failures. 0 regressions.**
|
||||
**1276 test functions. 50 test files. All green. 0 failures. 0 regressions.**
|
||||
|
||||
### Completed subsystems
|
||||
|
||||
@@ -415,6 +416,8 @@ simple doctrinal tick-exits (C11) ship first via T19 step 3; MALKHUT supersedes
|
||||
| **Zinc IPC** | `ipc/zinc_plane.py` | 8 | Real POSIX SHM, UVZINC01 framing |
|
||||
| **Control Plane** | `ipc/control_plane.py` | 3 | HOT_RELOAD_POLICY, kill switch |
|
||||
| **ClickHouse** | `storage/ch_store.py` | 9 | 5 tables, HTTP API |
|
||||
| **DuckDB Asset Store** | `storage/asset_store.py` | 30 | File-backed asset universe, query, sync, persistence |
|
||||
| **Asset Bridge** | `training/asset_bridge.py` | 49 | Directory ↔ classification sync |
|
||||
| **CMA-ES Training** | `training/cma_trainer.py` | 65 | Behavior-driven scenarios, auto-compile, label queries |
|
||||
| **Parallel Eval** | `training/parallel_eval.py` | 16 | 9x speedup, zero fidelity loss, ProcessPoolExecutor |
|
||||
| **Policy Registry** | `training/registry.py` | 14 | CANDIDATE → ACTIVE lifecycle |
|
||||
|
||||
@@ -1 +1,3 @@
|
||||
from malkhut.storage.ch_store import MalkhutCHStore
|
||||
from malkhut.storage.asset_store import AssetStore
|
||||
|
||||
__all__ = ["AssetStore"]
|
||||
|
||||
272
MALKHUT/malkhut/storage/asset_store.py
Normal file
272
MALKHUT/malkhut/storage/asset_store.py
Normal file
@@ -0,0 +1,272 @@
|
||||
"""
|
||||
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()
|
||||
253
MALKHUT/malkhut/tests/test_asset_store.py
Normal file
253
MALKHUT/malkhut/tests/test_asset_store.py
Normal file
@@ -0,0 +1,253 @@
|
||||
"""
|
||||
Tests for DuckDB asset store — schema, sync, queries, persistence.
|
||||
|
||||
Covers:
|
||||
- Schema creation and table structure
|
||||
- Sync from in-memory profiles → DuckDB
|
||||
- Asset CRUD (upsert, query, get)
|
||||
- Exchange-asset mapping
|
||||
- Behavior profiles
|
||||
- Cross-exchange queries
|
||||
- Blockchain filtering
|
||||
- Case-insensitive exchange lookup
|
||||
- Persistence across connections
|
||||
- Performance baseline
|
||||
"""
|
||||
import pytest
|
||||
import os
|
||||
from pathlib import Path
|
||||
|
||||
from malkhut.storage.asset_store import AssetStore
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def store(tmp_path):
|
||||
"""Create a fresh DuckDB store for each test."""
|
||||
db_path = str(tmp_path / "test_assets.duckdb")
|
||||
s = AssetStore(db_path=db_path)
|
||||
yield s
|
||||
s.close()
|
||||
|
||||
|
||||
class TestSchema:
|
||||
def test_tables_created(self, store):
|
||||
tables = store.conn.execute(
|
||||
"SELECT table_name FROM information_schema.tables WHERE table_schema='main'"
|
||||
).fetchall()
|
||||
names = {t[0] for t in tables}
|
||||
assert "assets" in names
|
||||
assert "exchanges" in names
|
||||
assert "asset_exchanges" in names
|
||||
assert "behavior_profiles" in names
|
||||
|
||||
def test_asset_table_columns(self, store):
|
||||
cols = store.conn.execute(
|
||||
"SELECT column_name FROM information_schema.columns WHERE table_name='assets'"
|
||||
).fetchall()
|
||||
col_names = [c[0] for c in cols]
|
||||
assert "symbol" in col_names
|
||||
assert "base_asset" in col_names
|
||||
assert "coingecko_id" in col_names
|
||||
assert "sectors" in col_names
|
||||
assert "tick_size" in col_names
|
||||
|
||||
def test_indexes_created(self, store):
|
||||
indexes = store.conn.execute(
|
||||
"SELECT name FROM sqlite_master WHERE type='index' AND tbl_name='assets'"
|
||||
).fetchall()
|
||||
idx_names = {i[0] for i in indexes}
|
||||
assert any("blockchain" in n for n in idx_names)
|
||||
assert any("coingecko" in n for n in idx_names)
|
||||
|
||||
|
||||
class TestSyncFromProfiles:
|
||||
def test_sync_returns_count(self, store):
|
||||
count = store.sync_from_profiles()
|
||||
assert count >= 13
|
||||
|
||||
def test_sync_populates_assets(self, store):
|
||||
store.sync_from_profiles()
|
||||
assert store.asset_count() >= 13
|
||||
|
||||
def test_sync_populates_exchanges(self, store):
|
||||
store.sync_from_profiles()
|
||||
assert store.exchange_count() >= 3
|
||||
|
||||
def test_sync_populates_asset_exchanges(self, store):
|
||||
store.sync_from_profiles()
|
||||
exs = store.get_asset_exchanges("BTCUSDT")
|
||||
assert len(exs) >= 1
|
||||
|
||||
def test_sync_populates_behaviors(self, store):
|
||||
store.sync_from_profiles()
|
||||
row = store.conn.execute(
|
||||
"SELECT * FROM behavior_profiles WHERE symbol = 'BTCUSDT'"
|
||||
).fetchone()
|
||||
assert row is not None
|
||||
|
||||
def test_sync_is_idempotent(self, store):
|
||||
store.sync_from_profiles()
|
||||
count1 = store.asset_count()
|
||||
store.sync_from_profiles()
|
||||
count2 = store.asset_count()
|
||||
assert count1 == count2
|
||||
|
||||
|
||||
class TestAssetCRUD:
|
||||
def test_get_asset_hit(self, store):
|
||||
store.sync_from_profiles()
|
||||
btc = store.get_asset("BTCUSDT")
|
||||
assert btc is not None
|
||||
assert btc["base_asset"] == "BTC"
|
||||
assert btc["name"] == "BTC"
|
||||
|
||||
def test_get_asset_miss(self, store):
|
||||
store.sync_from_profiles()
|
||||
assert store.get_asset("NONEXISTENT") is None
|
||||
|
||||
def test_get_asset_all_fields(self, store):
|
||||
store.sync_from_profiles()
|
||||
btc = store.get_asset("BTCUSDT")
|
||||
assert btc["coingecko_id"] == "bitcoin"
|
||||
assert btc["cmc_id"] == 1
|
||||
assert btc["blockchain"] == "bitcoin"
|
||||
assert btc["unified_symbol"] == "BTC/USDT"
|
||||
assert btc["quote_currency"] == "USDT"
|
||||
assert btc["tick_size"] == 0.1
|
||||
assert btc["lot_size"] == 0.001
|
||||
|
||||
def test_query_assets_by_blockchain(self, store):
|
||||
store.sync_from_profiles()
|
||||
eth_assets = store.query_assets(blockchain="ethereum")
|
||||
assert len(eth_assets) >= 5
|
||||
symbols = {a["symbol"] for a in eth_assets}
|
||||
assert "ETHUSDT" in symbols
|
||||
assert "UNIUSDT" in symbols
|
||||
|
||||
def test_query_assets_by_coingecko(self, store):
|
||||
store.sync_from_profiles()
|
||||
results = store.query_assets(coingecko_id="bitcoin")
|
||||
assert len(results) == 1
|
||||
assert results[0]["symbol"] == "BTCUSDT"
|
||||
|
||||
def test_query_assets_by_sector(self, store):
|
||||
store.sync_from_profiles()
|
||||
defi = store.query_assets(sectors=["defi"])
|
||||
assert len(defi) >= 2
|
||||
symbols = {a["symbol"] for a in defi}
|
||||
assert "UNIUSDT" in symbols
|
||||
|
||||
def test_asset_count(self, store):
|
||||
store.sync_from_profiles()
|
||||
assert store.asset_count() == 13
|
||||
|
||||
|
||||
class TestExchangeMapping:
|
||||
def test_symbols_for_exchange(self, store):
|
||||
store.sync_from_profiles()
|
||||
binance_assets = store.symbols_for_exchange("binance")
|
||||
assert len(binance_assets) >= 13
|
||||
|
||||
def test_symbols_for_exchange_case_insensitive(self, store):
|
||||
store.sync_from_profiles()
|
||||
upper = store.symbols_for_exchange("BINANCE")
|
||||
lower = store.symbols_for_exchange("binance")
|
||||
assert upper == lower
|
||||
|
||||
def test_symbols_for_unknown_exchange(self, store):
|
||||
store.sync_from_profiles()
|
||||
assert store.symbols_for_exchange("KRAKEN") == []
|
||||
|
||||
def test_assets_on_blockchain(self, store):
|
||||
store.sync_from_profiles()
|
||||
eth = store.assets_on_blockchain("ethereum")
|
||||
assert len(eth) >= 5
|
||||
|
||||
def test_assets_on_unknown_blockchain(self, store):
|
||||
store.sync_from_profiles()
|
||||
assert store.assets_on_blockchain("nonexistent") == []
|
||||
|
||||
|
||||
class TestBehaviorProfiles:
|
||||
def test_behavior_stored(self, store):
|
||||
store.sync_from_profiles()
|
||||
row = store.conn.execute(
|
||||
"SELECT depth_amplitude_usd, vol_annualized_normal FROM behavior_profiles WHERE symbol='BTCUSDT'"
|
||||
).fetchone()
|
||||
assert row[0] == 750_000.0
|
||||
assert row[1] == 35.0
|
||||
|
||||
def test_behavior_template(self, store):
|
||||
store.sync_from_profiles()
|
||||
row = store.conn.execute(
|
||||
"SELECT template_name FROM behavior_profiles WHERE symbol='BTCUSDT'"
|
||||
).fetchone()
|
||||
assert row[0] == "institutional_blue_chip"
|
||||
|
||||
def test_behavior_reference_price(self, store):
|
||||
store.sync_from_profiles()
|
||||
row = store.conn.execute(
|
||||
"SELECT reference_price FROM behavior_profiles WHERE symbol='BTCUSDT'"
|
||||
).fetchone()
|
||||
assert row[0] == 64000.0
|
||||
|
||||
def test_behavior_liquidation(self, store):
|
||||
store.sync_from_profiles()
|
||||
row = store.conn.execute(
|
||||
"SELECT liq_speed, liq_recovery FROM behavior_profiles WHERE symbol='BTCUSDT'"
|
||||
).fetchone()
|
||||
assert row[0] == "slow"
|
||||
assert row[1] == "fast"
|
||||
|
||||
|
||||
class TestPersistence:
|
||||
def test_data_survives_reopen(self, tmp_path):
|
||||
db_path = str(tmp_path / "persist.duckdb")
|
||||
s1 = AssetStore(db_path=db_path)
|
||||
s1.sync_from_profiles()
|
||||
s1.close()
|
||||
|
||||
s2 = AssetStore(db_path=db_path)
|
||||
assert s2.asset_count() >= 13
|
||||
btc = s2.get_asset("BTCUSDT")
|
||||
assert btc["base_asset"] == "BTC"
|
||||
s2.close()
|
||||
|
||||
def test_exchange_mapping_survives(self, tmp_path):
|
||||
db_path = str(tmp_path / "persist.duckdb")
|
||||
s1 = AssetStore(db_path=db_path)
|
||||
s1.sync_from_profiles()
|
||||
exs = s1.get_asset_exchanges("BTCUSDT")
|
||||
s1.close()
|
||||
|
||||
s2 = AssetStore(db_path=db_path)
|
||||
exs2 = s2.get_asset_exchanges("BTCUSDT")
|
||||
assert exs == exs2
|
||||
s2.close()
|
||||
|
||||
|
||||
class TestPerformanceBaseline:
|
||||
def test_sync_speed(self, store):
|
||||
import time
|
||||
t0 = time.time()
|
||||
store.sync_from_profiles()
|
||||
elapsed = time.time() - t0
|
||||
assert elapsed < 1.0, f"Sync took {elapsed:.2f}s"
|
||||
|
||||
def test_query_speed(self, store):
|
||||
import time
|
||||
store.sync_from_profiles()
|
||||
t0 = time.time()
|
||||
for _ in range(100):
|
||||
store.query_assets(blockchain="ethereum")
|
||||
elapsed = time.time() - t0
|
||||
assert elapsed < 1.0, f"100 queries took {elapsed:.2f}s"
|
||||
|
||||
def test_get_speed(self, store):
|
||||
import time
|
||||
store.sync_from_profiles()
|
||||
t0 = time.time()
|
||||
for _ in range(1000):
|
||||
store.get_asset("BTCUSDT")
|
||||
elapsed = time.time() - t0
|
||||
assert elapsed < 1.0, f"1000 gets took {elapsed:.2f}s"
|
||||
Reference in New Issue
Block a user