malkhut: asset-faithful book generation with composable toggles
Three independently toggleable features: 1. Asset-faithful depth/spread: levels sized by OB study power-law per asset 2. Intraday volume clock: depth scales by time-of-day (peak/trough) 3. Realistic spread: per-asset spread from OB study + Flight7 Composable via BookGenerationConfig toggles: use_asset_faithful_depth, use_asset_faithful_spread, use_intraday_clock, use_weekend_mode, use_stress_mode, use_fragility, worst_case_mode worst_case_mode overrides everything for max adversarial learning: spread * stress_mult, depth * fragility, no intraday/weekend. DuckDB registry for online updates: AssetRegistry: upsert/get/list/delete/query RuntimeProfileCache: hot-reload during CWM runs upsert_from_csv/export_csv: pipeline support upsert_all_from_asset_behaviors(): seed from OB study Results (8 assets): BTC: spread 0.031 bps, depth $350M (normal) / $4.9M (worst) DOGE: spread 2.86 bps, depth $2M (normal) / $132K (worst) ADA: spread 11.8 bps, depth $10M (normal) / $511K (worst) Intraday: BTC peak/trough = 2.8x depth ratio All 99 tests green (31 new + 68 existing).
This commit is contained in:
168
MALKHUT/malkhut/training/asset_registry.py
Normal file
168
MALKHUT/malkhut/training/asset_registry.py
Normal file
@@ -0,0 +1,168 @@
|
||||
"""
|
||||
Asset Book Profile Registry — DuckDB persistence + online update tooling.
|
||||
|
||||
Provides upsert/query for per-asset book generation profiles.
|
||||
Profiles can be updated:
|
||||
1. One-shot: upsert_all_from_asset_behaviors() seeds all 13 assets
|
||||
2. Online: upsert_profile(symbol, ...) updates a single asset
|
||||
3. Pipeline: upsert_from_csv(path) bulk-loads from a CSV
|
||||
4. Runtime override: RuntimeProfileCache for hot-reload during CWM runs
|
||||
|
||||
Usage:
|
||||
from malkhut.training.asset_registry import AssetRegistry
|
||||
reg = AssetRegistry()
|
||||
reg.upsert_all_from_asset_behaviors()
|
||||
profile = reg.get_profile("BTCUSDT")
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import csv
|
||||
import os
|
||||
from typing import Dict, List, Optional
|
||||
|
||||
from malkhut.training.asset_book_profile import AssetBookProfile
|
||||
|
||||
_DEFAULT_DB = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))),
|
||||
"data", "asset_registry.db")
|
||||
|
||||
_CREATE_SQL = """
|
||||
CREATE TABLE IF NOT EXISTS asset_book_profiles (
|
||||
symbol TEXT PRIMARY KEY,
|
||||
depth_amplitude_usd DOUBLE, depth_alpha DOUBLE, depth_fragility DOUBLE,
|
||||
depth_at_10bps_usd DOUBLE, depth_at_100bps_usd DOUBLE,
|
||||
spread_normal_bps DOUBLE, spread_stress_mult DOUBLE,
|
||||
flow_orders_per_sec DOUBLE, flow_cancel_fill_ratio DOUBLE,
|
||||
flow_median_order_usd DOUBLE, flow_avg_trade_usd DOUBLE,
|
||||
vol_annualized_normal DOUBLE, vol_annualized_crisis DOUBLE,
|
||||
vol_garch_alpha DOUBLE, vol_garch_beta DOUBLE, vol_half_life_hours DOUBLE,
|
||||
intraday_peak_hour_utc INTEGER, intraday_trough_hour_utc INTEGER,
|
||||
intraday_ratio DOUBLE,
|
||||
weekend_vol_mult DOUBLE, weekend_volume_mult DOUBLE, weekend_spread_mult DOUBLE,
|
||||
mm_max_inventory_usd DOUBLE, mm_pull_speed_ms DOUBLE, mm_margin_bps DOUBLE,
|
||||
avg_level_size_usd DOUBLE, typical_num_levels INTEGER, reference_price DOUBLE,
|
||||
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
|
||||
)
|
||||
"""
|
||||
|
||||
|
||||
class AssetRegistry:
|
||||
"""DuckDB-backed asset profile registry with online update support."""
|
||||
|
||||
def __init__(self, db_path: str = _DEFAULT_DB) -> None:
|
||||
os.makedirs(os.path.dirname(db_path), exist_ok=True)
|
||||
import duckdb
|
||||
self._db_path = db_path
|
||||
self._conn = duckdb.connect(db_path)
|
||||
self._conn.execute(_CREATE_SQL)
|
||||
|
||||
def upsert_profile(self, profile: AssetBookProfile) -> None:
|
||||
d = profile.to_dict()
|
||||
cols = list(d.keys())
|
||||
placeholders = ", ".join(["?"] * len(cols))
|
||||
col_str = ", ".join(cols)
|
||||
self._conn.execute(
|
||||
f"INSERT INTO asset_book_profiles ({col_str}) VALUES ({placeholders}) "
|
||||
f"ON CONFLICT (symbol) DO UPDATE SET {', '.join(f'{c}=excluded.{c}' for c in cols)}",
|
||||
list(d.values()),
|
||||
)
|
||||
|
||||
def get_profile(self, symbol: str) -> Optional[AssetBookProfile]:
|
||||
rows = self._conn.execute(
|
||||
"SELECT * FROM asset_book_profiles WHERE symbol = ?", [symbol]
|
||||
).fetchall()
|
||||
if not rows:
|
||||
return None
|
||||
cols = [desc[0] for desc in self._conn.description]
|
||||
return AssetBookProfile.from_dict(dict(zip(cols, rows[0])))
|
||||
|
||||
def list_profiles(self) -> List[AssetBookProfile]:
|
||||
rows = self._conn.execute("SELECT * FROM asset_book_profiles").fetchall()
|
||||
cols = [desc[0] for desc in self._conn.description]
|
||||
return [AssetBookProfile.from_dict(dict(zip(cols, r))) for r in rows]
|
||||
|
||||
def list_symbols(self) -> List[str]:
|
||||
rows = self._conn.execute("SELECT symbol FROM asset_book_profiles").fetchall()
|
||||
return [r[0] for r in rows]
|
||||
|
||||
def delete_profile(self, symbol: str) -> None:
|
||||
self._conn.execute("DELETE FROM asset_book_profiles WHERE symbol = ?", [symbol])
|
||||
|
||||
def upsert_all_from_asset_behaviors(self) -> int:
|
||||
from malkhut.training.asset_book_profile import build_profile_from_behavior
|
||||
from malkhut.training.asset_behavior import list_behavior_symbols
|
||||
count = 0
|
||||
for sym in list_behavior_symbols():
|
||||
try:
|
||||
profile = build_profile_from_behavior(sym)
|
||||
self.upsert_profile(profile)
|
||||
count += 1
|
||||
except Exception:
|
||||
continue
|
||||
return count
|
||||
|
||||
def upsert_from_csv(self, csv_path: str) -> int:
|
||||
count = 0
|
||||
with open(csv_path, "r") as f:
|
||||
reader = csv.DictReader(f)
|
||||
for row in reader:
|
||||
try:
|
||||
profile = AssetBookProfile.from_dict(
|
||||
{k: float(v) if k not in ("symbol",) else v
|
||||
for k, v in row.items() if hasattr(AssetBookProfile, k)}
|
||||
)
|
||||
self.upsert_profile(profile)
|
||||
count += 1
|
||||
except Exception:
|
||||
continue
|
||||
return count
|
||||
|
||||
def export_csv(self, csv_path: str) -> int:
|
||||
profiles = self.list_profiles()
|
||||
if not profiles:
|
||||
return 0
|
||||
cols = list(profiles[0].to_dict().keys())
|
||||
with open(csv_path, "w", newline="") as f:
|
||||
writer = csv.DictWriter(f, fieldnames=cols)
|
||||
writer.writeheader()
|
||||
for p in profiles:
|
||||
writer.writerow(p.to_dict())
|
||||
return len(profiles)
|
||||
|
||||
def close(self) -> None:
|
||||
self._conn.close()
|
||||
|
||||
|
||||
class RuntimeProfileCache:
|
||||
"""Hot-reloadable in-memory cache of AssetBookProfiles.
|
||||
|
||||
CWM uses this to pick up profile updates mid-run without restart.
|
||||
Supports polling (check for updates) and push (explicit update).
|
||||
"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self._cache: Dict[str, AssetBookProfile] = {}
|
||||
|
||||
def get(self, symbol: str) -> Optional[AssetBookProfile]:
|
||||
return self._cache.get(symbol)
|
||||
|
||||
def put(self, profile: AssetBookProfile) -> None:
|
||||
self._cache[profile.symbol] = profile
|
||||
|
||||
def put_all(self, profiles: List[AssetBookProfile]) -> None:
|
||||
for p in profiles:
|
||||
self._cache[p.symbol] = p
|
||||
|
||||
def load_from_registry(self, registry: AssetRegistry, symbols: Optional[List[str]] = None) -> int:
|
||||
if symbols is None:
|
||||
profiles = registry.list_profiles()
|
||||
else:
|
||||
profiles = [registry.get_profile(s) for s in symbols]
|
||||
profiles = [p for p in profiles if p is not None]
|
||||
self.put_all(profiles)
|
||||
return len(profiles)
|
||||
|
||||
def has(self, symbol: str) -> bool:
|
||||
return symbol in self._cache
|
||||
|
||||
def symbols(self) -> List[str]:
|
||||
return list(self._cache.keys())
|
||||
Reference in New Issue
Block a user