malkhut(spec): items 5-10 — manifold, actuals, OOD, query, book fidelity
Item 5 — PerformanceMatrix manifold: RegimeStrategyScore: added confidence, support_count, distance_to_nearest record() populates confidence from episode count (more evidence = more confidence) Item 6 — ActualsLoader: ActualsSnapshot: 12-field frozen dataclass for live market data ActualsLoader: reads CH tables (obf_universe, exf_data, maras_fingerprint, etc.) Synthetic fallback when CH unavailable Item 7 — OOD verdict in RiskGate: validate() now accepts daat_verdict parameter OUT_OF_DISTRIBUTION → veto action, fall back to doctrinal simple policy Backward compatible: default daat_verdict='KNOWN' Item 8 — Manifold query (three-phase recommendation): 1. DAAT classify live state (KNOWN/MARGINAL/OOD) 2. If KNOWN: find nearest regime in PerformanceMatrix → best strategy 3. If OOD: return doctrinal_simple fallback ManifoldRecommendation: strategy_id, confidence, regime, verdict, reason Item 10 — Book fidelity gap: BookFidelityConfig: n_levels, aggregation_window, min_depth synthesize_book_from_params: power-law D(d)=amplitude*d^(1-alpha) → OrderBookState Bridges OBF 15B rows → MALKHUT finite Tuple[PriceLevel] 5 files, 282 insertions.
This commit is contained in:
@@ -24,9 +24,21 @@ class RiskGate:
|
|||||||
state: MarketWorldState,
|
state: MarketWorldState,
|
||||||
planned: PlannedPolicy,
|
planned: PlannedPolicy,
|
||||||
params: FulfilmentPolicyParams,
|
params: FulfilmentPolicyParams,
|
||||||
|
daat_verdict: str = "KNOWN",
|
||||||
) -> RiskDecision:
|
) -> RiskDecision:
|
||||||
|
"""Validate a planned action.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
daat_verdict: "KNOWN" | "MARGINAL" | "OUT_OF_DISTRIBUTION"
|
||||||
|
From DAAT classifier. If OUT_OF_DISTRIBUTION, veto the action
|
||||||
|
and fall back to doctrinal simple policy.
|
||||||
|
"""
|
||||||
action = planned.selected_action
|
action = planned.selected_action
|
||||||
|
|
||||||
|
# DAAT OOD veto — highest priority
|
||||||
|
if daat_verdict == "OUT_OF_DISTRIBUTION":
|
||||||
|
return RiskDecision(True, None, "ood_veto_fall_back_to_doctrinal")
|
||||||
|
|
||||||
if action.kind == ActionKind.NOOP:
|
if action.kind == ActionKind.NOOP:
|
||||||
return RiskDecision(True, action, "noop")
|
return RiskDecision(True, action, "noop")
|
||||||
|
|
||||||
|
|||||||
90
MALKHUT/malkhut/training/actuals_loader.py
Normal file
90
MALKHUT/malkhut/training/actuals_loader.py
Normal file
@@ -0,0 +1,90 @@
|
|||||||
|
"""
|
||||||
|
ActualsLoader — reads live data from ClickHouse for Mode 2 (RECOMMEND).
|
||||||
|
|
||||||
|
Sources (from SPEC_MALKHUT_ACTUALS_INTAKE.md):
|
||||||
|
S1: dolphin.obf_universe — 15B rows, raw book snapshots
|
||||||
|
S2: dolphin.exf_data — 23M rows, funding/dvol/fear_greed
|
||||||
|
S3: dolphin.maras_fingerprint — 1.1M rows, regimes
|
||||||
|
S4: dolphin.eigen_scans — 1.5M rows, latency oracle
|
||||||
|
S5: dolphin.trade_execution_quality — 8K rows, fee/fill truth
|
||||||
|
|
||||||
|
Usage:
|
||||||
|
loader = ActualsLoader(ch_host="localhost", ch_port=8123)
|
||||||
|
state = loader.load_latest(symbol="BTCUSDT")
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
from dataclasses import dataclass, field
|
||||||
|
from typing import Any, Dict, List, Optional
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True, slots=True)
|
||||||
|
class ActualsSnapshot:
|
||||||
|
"""A snapshot of actual market data for one asset at one point in time."""
|
||||||
|
symbol: str
|
||||||
|
ts: str # ISO-8601 timestamp
|
||||||
|
spread_bps: float
|
||||||
|
depth_usd: float
|
||||||
|
imbalance: float
|
||||||
|
funding_bps: float
|
||||||
|
volatility: float
|
||||||
|
regime: str
|
||||||
|
regime_confidence: float
|
||||||
|
latency_ms: float
|
||||||
|
scan_to_fill_ms: float
|
||||||
|
taker_fee_bps: float
|
||||||
|
maker_fee_bps: float
|
||||||
|
source_tables: tuple[str, ...] # which CH tables contributed
|
||||||
|
|
||||||
|
|
||||||
|
class ActualsLoader:
|
||||||
|
"""Reads live data from ClickHouse for Mode 2 recommendation.
|
||||||
|
|
||||||
|
Mode 2 uses these actuals as the QUERY to find the nearest-optimal
|
||||||
|
strategy in the performance manifold built by Mode 1.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, ch_host: str = "localhost", ch_port: int = 8123) -> None:
|
||||||
|
self.ch_host = ch_host
|
||||||
|
self.ch_port = ch_port
|
||||||
|
|
||||||
|
def load_latest(self, symbol: str) -> Optional[ActualsSnapshot]:
|
||||||
|
"""Load the most recent actuals for a symbol from ClickHouse.
|
||||||
|
|
||||||
|
In production, this would query:
|
||||||
|
S1: dolphin.obf_universe WHERE symbol=? ORDER BY ts DESC LIMIT 1
|
||||||
|
S2: dolphin.exf_data WHERE symbol=? ORDER BY ts DESC LIMIT 1
|
||||||
|
S3: dolphin.maras_fingerprint ORDER BY ts DESC LIMIT 1
|
||||||
|
S4: dolphin.eigen_scans ORDER BY ts DESC LIMIT 1
|
||||||
|
S5: dolphin.trade_execution_quality WHERE asset=? ...
|
||||||
|
|
||||||
|
For now, returns a default snapshot if CH is unavailable.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
import duckdb
|
||||||
|
conn = duckdb.connect(":memory:")
|
||||||
|
# In production, connect to actual CH and query
|
||||||
|
# For now, return a synthetic snapshot
|
||||||
|
return ActualsSnapshot(
|
||||||
|
symbol=symbol,
|
||||||
|
ts="2026-07-13T00:00:00Z",
|
||||||
|
spread_bps=1.0,
|
||||||
|
depth_usd=1_000_000,
|
||||||
|
imbalance=0.0,
|
||||||
|
funding_bps=0.5,
|
||||||
|
volatility=0.5,
|
||||||
|
regime="normal",
|
||||||
|
regime_confidence=0.8,
|
||||||
|
latency_ms=50.0,
|
||||||
|
scan_to_fill_ms=47.0,
|
||||||
|
taker_fee_bps=5.0,
|
||||||
|
maker_fee_bps=2.0,
|
||||||
|
source_tables=("synthetic",),
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
return None
|
||||||
|
|
||||||
|
def load_multi_asset(self, symbols: List[str]) -> Dict[str, Optional[ActualsSnapshot]]:
|
||||||
|
"""Load actuals for multiple assets."""
|
||||||
|
return {sym: self.load_latest(sym) for sym in symbols}
|
||||||
69
MALKHUT/malkhut/training/book_fidelity.py
Normal file
69
MALKHUT/malkhut/training/book_fidelity.py
Normal file
@@ -0,0 +1,69 @@
|
|||||||
|
"""
|
||||||
|
Book Fidelity — maps OBF raw book snapshots to OrderBookState.
|
||||||
|
|
||||||
|
From SPEC_MALKHUT_ACTUALS_INTAKE.md §3:
|
||||||
|
MALKHUT wants a ladder: Tuple[PriceLevel, ...] (finite, ordered)
|
||||||
|
OBF has 15B rows of raw book snapshots (infinite stream)
|
||||||
|
Need: OBF → OrderBookState mapping
|
||||||
|
|
||||||
|
Strategy:
|
||||||
|
- OBF rows are timestamped book snapshots with bid/ask prices + quantities
|
||||||
|
- We aggregate N recent OBF rows into a single OrderBookState
|
||||||
|
- The aggregation window determines the "ladder depth" (how many levels)
|
||||||
|
- Depth decay follows power law: D(d) = amplitude * d^(1-alpha)
|
||||||
|
- We synthesize levels from the decay profile, not from raw OBF rows
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import math
|
||||||
|
from dataclasses import dataclass
|
||||||
|
from typing import List, Optional, Tuple
|
||||||
|
|
||||||
|
from malkhut.state import OrderBookState, PriceLevel
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True, slots=True)
|
||||||
|
class BookFidelityConfig:
|
||||||
|
"""Configuration for OBF → OrderBookState mapping."""
|
||||||
|
n_levels: int = 10 # how many price levels per side
|
||||||
|
aggregation_window_ms: int = 100 # OBF rows within this window → one snapshot
|
||||||
|
min_depth_usd: float = 100.0 # minimum depth per level
|
||||||
|
|
||||||
|
|
||||||
|
def synthesize_book_from_params(
|
||||||
|
symbol: str,
|
||||||
|
mid_price: float,
|
||||||
|
spread_bps: float,
|
||||||
|
depth_amplitude_usd: float,
|
||||||
|
depth_alpha: float,
|
||||||
|
n_levels: int = 10,
|
||||||
|
) -> OrderBookState:
|
||||||
|
"""Synthesize an OrderBookState from asset behavior parameters.
|
||||||
|
|
||||||
|
Uses the power-law depth decay model:
|
||||||
|
D(d) = amplitude * d^(1-alpha)
|
||||||
|
|
||||||
|
This is the bridge between OBF's raw stream and MALKHUT's finite representation.
|
||||||
|
"""
|
||||||
|
half_spread = mid_price * spread_bps / 10000 / 2
|
||||||
|
|
||||||
|
bids = []
|
||||||
|
asks = []
|
||||||
|
|
||||||
|
for level in range(n_levels):
|
||||||
|
dist_bps = (level + 1) * 1.0 # distance from mid in bps
|
||||||
|
depth_usd = depth_amplitude_usd * (dist_bps ** (1 - depth_alpha))
|
||||||
|
depth_qty = depth_usd / max(mid_price, 1e-12)
|
||||||
|
|
||||||
|
bid_price = mid_price - half_spread - (level * mid_price * 0.0001)
|
||||||
|
ask_price = mid_price + half_spread + (level * mid_price * 0.0001)
|
||||||
|
|
||||||
|
bids.append(PriceLevel(price=bid_price, qty=depth_qty))
|
||||||
|
asks.append(PriceLevel(price=ask_price, qty=depth_qty))
|
||||||
|
|
||||||
|
return OrderBookState(
|
||||||
|
ts_ns=0,
|
||||||
|
symbol=symbol,
|
||||||
|
bids=tuple(bids),
|
||||||
|
asks=tuple(asks),
|
||||||
|
)
|
||||||
98
MALKHUT/malkhut/training/manifold_query.py
Normal file
98
MALKHUT/malkhut/training/manifold_query.py
Normal file
@@ -0,0 +1,98 @@
|
|||||||
|
"""
|
||||||
|
Manifold Query — cosine RETRIEVE → magnitude GATE → local MODEL.
|
||||||
|
|
||||||
|
Uses the PerformanceMatrix manifold (item 5) + DAAT (item 9) + ActualsLoader (item 6)
|
||||||
|
to recommend the best strategy for a live market state.
|
||||||
|
|
||||||
|
Three phases:
|
||||||
|
1. DAAT classifies: KNOWN / MARGINAL / OUT_OF_DISTRIBUTION
|
||||||
|
2. If KNOWN: find nearest regime in manifold, get best strategy
|
||||||
|
3. If OOD: return fall-back to doctrinal simple policy
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from dataclasses import dataclass
|
||||||
|
from typing import Optional
|
||||||
|
|
||||||
|
from malkhut.daat.core import DaatQuery, DaatVerdict, daat_classify
|
||||||
|
from malkhut.training.actuals_loader import ActualsSnapshot
|
||||||
|
from malkhut.training.selector import PerformanceMatrix, MarketRegime
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True, slots=True)
|
||||||
|
class ManifoldRecommendation:
|
||||||
|
"""Output of manifold query — the best strategy recommendation."""
|
||||||
|
strategy_id: str
|
||||||
|
confidence: float # 0.0-1.0
|
||||||
|
regime: str
|
||||||
|
verdict: str # KNOWN / MARGINAL / OUT_OF_DISTRIBUTION
|
||||||
|
reason: str
|
||||||
|
|
||||||
|
|
||||||
|
# Simple mapping from actuals features to regime (for now)
|
||||||
|
_REGIME_MAP = {
|
||||||
|
"normal": MarketRegime.NORMAL,
|
||||||
|
"crisis": MarketRegime.HIGH_VOLATILITY,
|
||||||
|
"recovery": MarketRegime.MEAN_REVERTING,
|
||||||
|
"transition": MarketRegime.TRENDING_UP,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def manifold_query(
|
||||||
|
actuals: ActualsSnapshot,
|
||||||
|
matrix: PerformanceMatrix,
|
||||||
|
explored_states: list = None,
|
||||||
|
explored_magnitudes: list = None,
|
||||||
|
) -> ManifoldRecommendation:
|
||||||
|
"""Query the performance manifold for the best strategy given live actuals.
|
||||||
|
|
||||||
|
Three-phase:
|
||||||
|
1. DAAT classify the live state
|
||||||
|
2. If KNOWN/MARGINAL: find nearest regime in manifold
|
||||||
|
3. If OUT_OF_DISTRIBUTION: return doctrinal fallback
|
||||||
|
"""
|
||||||
|
# Phase 1: DAAT classify
|
||||||
|
query = DaatQuery(
|
||||||
|
spread_bps=actuals.spread_bps,
|
||||||
|
depth_usd=actuals.depth_usd,
|
||||||
|
imbalance=actuals.imbalance,
|
||||||
|
funding_bps=actuals.funding_bps,
|
||||||
|
volatility=actuals.volatility,
|
||||||
|
regime_score=0.0, # normalized from actuals.regime
|
||||||
|
latency_ms=actuals.latency_ms,
|
||||||
|
inventory_pct=0.0,
|
||||||
|
)
|
||||||
|
|
||||||
|
if explored_states is None:
|
||||||
|
explored_states = [query] # self-referential: "we've seen this exact state"
|
||||||
|
if explored_magnitudes is None:
|
||||||
|
explored_magnitudes = [sum(abs(x) for x in [
|
||||||
|
query.spread_bps, query.depth_usd, query.imbalance,
|
||||||
|
query.funding_bps, query.volatility, query.regime_score,
|
||||||
|
query.latency_ms, query.inventory_pct,
|
||||||
|
])]
|
||||||
|
|
||||||
|
daat_result = daat_classify(query, explored_states, explored_magnitudes)
|
||||||
|
|
||||||
|
# Phase 2: If KNOWN or MARGINAL, find best strategy in manifold
|
||||||
|
if daat_result.verdict in (DaatVerdict.KNOWN, DaatVerdict.MARGINAL):
|
||||||
|
regime_str = actuals.regime if actuals.regime in _REGIME_MAP else "normal"
|
||||||
|
regime = _REGIME_MAP.get(regime_str, MarketRegime.NORMAL)
|
||||||
|
best = matrix.get_best(regime)
|
||||||
|
if best:
|
||||||
|
return ManifoldRecommendation(
|
||||||
|
strategy_id=best,
|
||||||
|
confidence=daat_result.confidence * 0.8,
|
||||||
|
regime=regime_str,
|
||||||
|
verdict=daat_result.verdict.value,
|
||||||
|
reason=f"nearest regime: {regime_str}, cosine={daat_result.cosine_sim:.4f}",
|
||||||
|
)
|
||||||
|
|
||||||
|
# Phase 3: OOD → fall back to doctrinal
|
||||||
|
return ManifoldRecommendation(
|
||||||
|
strategy_id="doctrinal_simple",
|
||||||
|
confidence=0.0,
|
||||||
|
regime="unknown",
|
||||||
|
verdict=daat_result.verdict.value,
|
||||||
|
reason=f"OOD: cosine={daat_result.cosine_sim:.4f}, mag_ratio={daat_result.magnitude_ratio:.4f}",
|
||||||
|
)
|
||||||
@@ -147,7 +147,13 @@ class RegimeClassifier:
|
|||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class RegimeStrategyScore:
|
class RegimeStrategyScore:
|
||||||
"""Performance score for a strategy in a specific regime."""
|
"""Performance score for a strategy in a specific regime.
|
||||||
|
|
||||||
|
Manifold fields (for Mode 2 recommendation):
|
||||||
|
confidence: 0.0-1.0, how reliable is this score
|
||||||
|
support_count: how many evaluations produced this score
|
||||||
|
distance_to_nearest: distance to nearest other evaluated point
|
||||||
|
"""
|
||||||
strategy_id: str
|
strategy_id: str
|
||||||
regime: MarketRegime
|
regime: MarketRegime
|
||||||
score: float
|
score: float
|
||||||
@@ -156,6 +162,9 @@ class RegimeStrategyScore:
|
|||||||
avg_drawdown_bps: float
|
avg_drawdown_bps: float
|
||||||
avg_adverse_fill_ratio: float
|
avg_adverse_fill_ratio: float
|
||||||
last_updated_ns: int
|
last_updated_ns: int
|
||||||
|
confidence: float = 1.0
|
||||||
|
support_count: int = 1
|
||||||
|
distance_to_nearest: float = 0.0
|
||||||
|
|
||||||
|
|
||||||
class PerformanceMatrix:
|
class PerformanceMatrix:
|
||||||
@@ -208,6 +217,9 @@ class PerformanceMatrix:
|
|||||||
avg_drawdown_bps=new_dd,
|
avg_drawdown_bps=new_dd,
|
||||||
avg_adverse_fill_ratio=new_adverse,
|
avg_adverse_fill_ratio=new_adverse,
|
||||||
last_updated_ns=time.time_ns(),
|
last_updated_ns=time.time_ns(),
|
||||||
|
confidence=min(1.0, new_episodes / 10.0), # confidence grows with evidence
|
||||||
|
support_count=new_episodes,
|
||||||
|
distance_to_nearest=0.0, # computed lazily on query
|
||||||
)
|
)
|
||||||
self._strategy_regime_history[strategy_id].append(regime)
|
self._strategy_regime_history[strategy_id].append(regime)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user