diff --git a/MALKHUT/malkhut/cwm/core.py b/MALKHUT/malkhut/cwm/core.py index d78e5f7..adb4acb 100644 --- a/MALKHUT/malkhut/cwm/core.py +++ b/MALKHUT/malkhut/cwm/core.py @@ -27,6 +27,7 @@ import numpy as np from malkhut.state import ( AccountState, FulfilmentPolicyParams, + FillQuality, MarketWorldState, Mode, OpenOrderState, @@ -564,6 +565,51 @@ class MinimalCryptoLOBCWM: positions=new_positions, ) + # ── Fill Quality computation ──────────────────────────────────────── + mid = state.book.mid if state.book.bids and state.book.asks else 0.0 + slippage_bps = 0.0 + if new_fill_qty > 0 and mid > 0 and new_fill_price > 0: + slippage_bps = abs(new_fill_price - mid) / mid * 10_000 + is_maker_fill = (our_action.order_type and our_action.order_type.value == "LIMIT") or our_action.post_only if isinstance(our_action, FulfilmentAction) else False + price_improvement_bps = 0.0 + if new_fill_qty > 0 and isinstance(our_action, FulfilmentAction) and our_action.post_only and our_action.side: + if our_action.side == Side.BUY and state.book.bids: + price_improvement_bps = (state.book.best_bid - new_fill_price) / max(state.book.best_bid, 1e-12) * 10_000 + elif our_action.side == Side.SELL and state.book.asks: + price_improvement_bps = (new_fill_price - state.book.best_ask) / max(state.book.best_ask, 1e-12) * 10_000 + post_fill_adverse = 0.0 + new_mid = book.mid if book.bids and book.asks else 0.0 + if new_fill_qty > 0 and mid > 0 and new_mid > 0: + if isinstance(our_action, FulfilmentAction) and our_action.side == Side.BUY: + post_fill_adverse = (new_mid - mid) / mid * 10_000 + elif isinstance(our_action, FulfilmentAction) and our_action.side == Side.SELL: + post_fill_adverse = (mid - new_mid) / mid * 10_000 + prev_fq = state.fill_quality + rolling_fill_rate = 0.0 + if prev_fq and prev_fq.filled: + rolling_fill_rate = 0.8 * prev_fq.rolling_fill_rate + 0.2 * (1.0 if new_fill_qty > 0 else 0.0) + elif new_fill_qty > 0: + rolling_fill_rate = 0.2 + spread_bps = book.spread_bps if book.bids and book.asks else 0.0 + fill_value = 0.0 + if new_fill_qty > 0: + quality = price_improvement_bps if is_maker_fill else max(0.0, spread_bps - slippage_bps) + fill_value = quality - abs(post_fill_adverse) * 0.5 + + fq = FillQuality( + filled=new_fill_qty > 0, + fill_qty=new_fill_qty, + fill_price=new_fill_price, + requested_qty=our_action.qty_fraction * state.account.available_balance / max(mid, 1e-12) if isinstance(our_action, FulfilmentAction) and our_action.qty_fraction > 0 and mid > 0 else 0.0, + slippage_bps=slippage_bps, + price_improvement_bps=price_improvement_bps, + levels_consumed=0, + is_maker_fill=is_maker_fill, + rolling_fill_rate=rolling_fill_rate, + post_fill_adverse_bps=post_fill_adverse, + fill_value_score=fill_value, + ) + return MarketWorldState( ts_ns=now_ts, mode=state.mode, @@ -577,8 +623,7 @@ class MinimalCryptoLOBCWM: volatility_state=state.volatility_state, market_regime=state.market_regime, feed_latency_ms=state.feed_latency_ms, - order_latency_ms=state.order_latency_ms, - rng_seed=state.rng_seed, + fill_quality=fq, ) def reward( diff --git a/MALKHUT/malkhut/cwm/hft_cwm.py b/MALKHUT/malkhut/cwm/hft_cwm.py index 6a5e1c0..65cd409 100644 --- a/MALKHUT/malkhut/cwm/hft_cwm.py +++ b/MALKHUT/malkhut/cwm/hft_cwm.py @@ -23,6 +23,7 @@ import numpy as np from malkhut.state import ( AccountState, FulfilmentPolicyParams, + FillQuality, MarketWorldState, OpenOrderState, OrderBookState, @@ -429,6 +430,17 @@ class HftBacktestCWM: positions=new_positions, ) + # ── Fill Quality computation (CORE metric) ────────────────────────── + fq = self._compute_fill_quality( + prev_state=state, + action=our_action, + new_fill_qty=new_fill_qty, + new_fill_price=new_fill_price, + book=book, + prev_book=state.book, + now_ts=now_ts, + ) + return MarketWorldState( ts_ns=now_ts, mode=state.mode, @@ -442,6 +454,95 @@ class HftBacktestCWM: volatility_state=state.volatility_state, market_regime=state.market_regime, feed_latency_ms=state.feed_latency_ms, + fill_quality=fq, + ) + + def _compute_fill_quality( + self, + prev_state: MarketWorldState, + action: FulfilmentAction, + new_fill_qty: float, + new_fill_price: float, + book: OrderBookState, + prev_book: OrderBookState, + now_ts: int, + ) -> FillQuality: + """Compute fill quality metrics for this transition. + + Fill quality is the CORE optimization target of MALKHUT. + Metrics: + - slippage_bps: how far from mid did we fill (aggressive) + - price_improvement_bps: how much better than touch (passive) + - levels_consumed: queue depth of fill + - is_maker_fill: passive vs aggressive + - rolling_fill_rate: recent fill success rate + - post_fill_adverse_bps: price movement after fill + - fill_value_score: composite optimization metric + """ + filled = new_fill_qty > 0 + mid = prev_book.mid if prev_book.bids and prev_book.asks else 0.0 + spread_bps = prev_book.spread_bps if prev_book.bids and prev_book.asks else 0.0 + + # Slippage: how far from mid did we fill? + slippage_bps = 0.0 + if filled and mid > 0 and new_fill_price > 0: + slippage_bps = abs(new_fill_price - mid) / mid * 10_000 + + # Price improvement: how much better than best bid/ask? + price_improvement_bps = 0.0 + if filled and action.post_only and action.side: + if action.side == Side.BUY and prev_book.bids: + price_improvement_bps = (prev_book.best_bid - new_fill_price) / max(prev_book.best_bid, 1e-12) * 10_000 + elif action.side == Side.SELL and prev_book.asks: + price_improvement_bps = (new_fill_price - prev_book.best_ask) / max(prev_book.best_ask, 1e-12) * 10_000 + + # Is maker fill? + is_maker = (action.order_type and action.order_type.value == "LIMIT") or action.post_only + + # Levels consumed (estimate: fill_qty / avg level qty) + levels_consumed = 0 + if filled and is_maker: + avg_level_qty = sum(l.qty for l in prev_book.asks if prev_book.asks) / max(len(prev_book.asks), 1) if action.side == Side.BUY else \ + sum(l.qty for l in prev_book.bids if prev_book.bids) / max(len(prev_book.bids), 1) + levels_consumed = max(1, int(new_fill_qty / max(avg_level_qty, 1e-12))) + + # Post-fill adverse: did price move against us? + post_fill_adverse = 0.0 + new_mid = book.mid if book.bids and book.asks else 0.0 + if filled and mid > 0 and new_mid > 0: + if action.side == Side.BUY: + post_fill_adverse = (new_mid - mid) / mid * 10_000 # negative = adverse + elif action.side == Side.SELL: + post_fill_adverse = (mid - new_mid) / mid * 10_000 # negative = adverse + + # Rolling fill rate (from state history) + prev_fq = prev_state.fill_quality + rolling_fill_rate = 0.0 + if prev_fq and prev_fq.filled: + rolling_fill_rate = 0.8 * prev_fq.rolling_fill_rate + 0.2 * (1.0 if filled else 0.0) + elif filled: + rolling_fill_rate = 0.2 + else: + rolling_fill_rate = 0.0 + + # Composite fill value score + fill_value = 0.0 + if filled: + quality = price_improvement_bps if is_maker else max(0.0, spread_bps - slippage_bps) + fill_value = quality - abs(post_fill_adverse) * 0.5 + + return FillQuality( + filled=filled, + fill_qty=new_fill_qty, + fill_price=new_fill_price, + requested_qty=action.qty_fraction * prev_state.account.available_balance / max(mid, 1e-12) if action.qty_fraction > 0 and mid > 0 else 0.0, + slippage_bps=slippage_bps, + price_improvement_bps=price_improvement_bps, + levels_consumed=levels_consumed, + is_maker_fill=is_maker, + rolling_fill_rate=rolling_fill_rate, + post_fill_adverse_bps=post_fill_adverse, + fill_value_score=fill_value, ) def reward( @@ -451,7 +552,18 @@ class HftBacktestCWM: next_state: MarketWorldState, params: FulfilmentPolicyParams, ) -> float: - """Same reward function as MinimalCryptoLOBCWM.""" + """Reward function — fill quality is the PRIMARY optimization target. + + MALKHUT is an execution improvement engine. Fill quality IS the core aim. + Reward = w_fill_probability * fill_value_score (PRIMARY) + + w_expected_pnl * pnl (secondary) + - w_adverse_selection * toxicity + - w_inventory_risk * inventory_risk + - w_tail_loss * tail_risk + - w_time_decay * time_in_loss + + w_fee_quality * maker_fee_benefit + - spread_cost - taker_fee + """ try: from malkhut.cwm.numba_core import compute_reward_vectorized @@ -469,7 +581,7 @@ class HftBacktestCWM: is_cross = action.kind.value == "CROSS_SPREAD" is_cancel = action.kind.value in ("CANCEL", "CANCEL_REPLACE") - return compute_reward_vectorized( + base_reward = compute_reward_vectorized( pnl, toxicity, churn, time_in_loss, spread_bps, inv_risk, tail_risk, params.w_expected_pnl, params.w_adverse_selection, @@ -481,36 +593,49 @@ class HftBacktestCWM: params.w_queue_priority, params.w_adverse_selection, ) except ImportError: - pass + # Fallback: Python path + fv = self.feature_extractor.extract(next_state).values + pnl = fv.get("pnl_bps", 0.0) + toxicity = fv.get("orderflow_toxicity", 0.0) + churn = fv.get("queue_churn_score", 0.0) + time_in_loss = fv.get("time_in_loss_s", 0.0) + spread_bps = fv.get("spread_bps", 0.0) - # Fallback: Python path - fv = self.feature_extractor.extract(next_state).values - pnl = fv.get("pnl_bps", 0.0) - toxicity = fv.get("orderflow_toxicity", 0.0) - churn = fv.get("queue_churn_score", 0.0) - time_in_loss = fv.get("time_in_loss_s", 0.0) - spread_bps = fv.get("spread_bps", 0.0) + base_reward = 0.0 + base_reward += params.w_expected_pnl * pnl + base_reward -= params.w_adverse_selection * toxicity + base_reward -= params.w_inventory_risk * self._inventory_risk(next_state) + base_reward -= params.w_tail_loss * self._tail_risk_proxy(next_state) + base_reward -= params.w_time_decay * math.log1p(max(time_in_loss, 0.0)) - reward = 0.0 - reward += params.w_expected_pnl * pnl - reward -= params.w_adverse_selection * toxicity - reward -= params.w_inventory_risk * self._inventory_risk(next_state) - reward -= params.w_tail_loss * self._tail_risk_proxy(next_state) - reward -= params.w_time_decay * math.log1p(max(time_in_loss, 0.0)) + if (action.order_type and action.order_type.value == "LIMIT") or action.post_only: + base_reward += params.w_fee_quality * max(0.0, -prev_state.venue.maker_fee_bps) - if (action.order_type and action.order_type.value == "LIMIT") or action.post_only: - reward += params.w_fee_quality * max(0.0, -prev_state.venue.maker_fee_bps) + if action.kind.value == "CROSS_SPREAD": + base_reward -= spread_bps + max(prev_state.venue.taker_fee_bps, 0.0) - if action.kind.value == "CROSS_SPREAD": - reward -= spread_bps + max(prev_state.venue.taker_fee_bps, 0.0) + if action.kind.value in ("CANCEL", "CANCEL_REPLACE"): + if toxicity > params.adverse_toxicity_cancel_threshold: + base_reward += params.w_adverse_selection * toxicity + if churn > params.queue_churn_cancel_threshold: + base_reward += params.w_queue_priority * churn - if action.kind.value in ("CANCEL", "CANCEL_REPLACE"): - if toxicity > params.adverse_toxicity_cancel_threshold: - reward += params.w_adverse_selection * toxicity - if churn > params.queue_churn_cancel_threshold: - reward += params.w_queue_priority * churn + # ── FILL QUALITY: the CORE reward signal ────────────────────────── + fq = next_state.fill_quality + fill_quality_reward = 0.0 + if fq: + # Primary: fill value score (price quality + fill success) + fill_quality_reward += params.w_fill_probability * fq.fill_value_score - return reward + # Bonus for maker fills that improve price + if fq.is_maker_fill and fq.price_improvement_bps > 0: + fill_quality_reward += params.w_fill_probability * fq.price_improvement_bps * 0.5 + + # Penalty for adverse selection after fill + if fq.filled and fq.post_fill_adverse_bps < 0: + fill_quality_reward += params.w_adverse_selection * fq.post_fill_adverse_bps + + return base_reward + fill_quality_reward def terminal(self, state: MarketWorldState, depth: int) -> bool: return depth <= 0 diff --git a/MALKHUT/malkhut/state.py b/MALKHUT/malkhut/state.py index 53de9b1..63714b5 100644 --- a/MALKHUT/malkhut/state.py +++ b/MALKHUT/malkhut/state.py @@ -267,6 +267,42 @@ class ExecutionIntent: reason: str +@dataclass(frozen=True, slots=True) +class FillQuality: + """Fill quality metrics — the CORE optimization target of MALKHUT. + + MALKHUT is an execution improvement engine. Fill quality IS the primary aim. + Every transition records these metrics. The reward function weights them heavily. + The PerformanceMatrix tracks them per (regime, strategy, venue). + """ + filled: bool = False + fill_qty: float = 0.0 + fill_price: float = 0.0 + requested_qty: float = 0.0 + + # How close to mid did we fill? (for aggressive: positive = slipped) + slippage_bps: float = 0.0 + + # For passive fills: how much better than best bid/ask? (positive = improvement) + price_improvement_bps: float = 0.0 + + # How many levels deep was the fill? + levels_consumed: int = 0 + + # Was this a maker (passive) or taker (aggressive) fill? + is_maker_fill: bool = False + + # Fill rate rolling window (updated each transition) + rolling_fill_rate: float = 0.0 + + # Adverse selection: price movement after fill (negative = adverse) + post_fill_adverse_bps: float = 0.0 + + # Fill value score: composite metric for optimization + # = fill_rate * price_quality - adverse_selection - slippage + fill_value_score: float = 0.0 + + @dataclass(frozen=True, slots=True) class MarketWorldState: """Complete CWM root state. Immutable for safe tree search.""" @@ -287,6 +323,8 @@ class MarketWorldState: order_latency_ms: float = 0.0 rng_seed: int = 0 + fill_quality: Optional[FillQuality] = None + @dataclass(frozen=True, slots=True) class FulfilmentPolicyParams: diff --git a/MALKHUT/malkhut/training/cma_trainer.py b/MALKHUT/malkhut/training/cma_trainer.py index 8dd4566..f4782ac 100644 --- a/MALKHUT/malkhut/training/cma_trainer.py +++ b/MALKHUT/malkhut/training/cma_trainer.py @@ -277,6 +277,10 @@ class EpisodeResult: final_equity: float = 0.0 max_position_qty: float = 0.0 diagnostics: Mapping[str, Any] = field(default_factory=dict) + # Fill quality (PRIMARY metrics) + avg_fill_value_score: float = 0.0 + avg_price_improvement_bps: float = 0.0 + avg_post_fill_adverse_bps: float = 0.0 # ============================================================================== @@ -1056,6 +1060,10 @@ class PolicyEvaluator: drawdown_bps=result.max_drawdown_bps, adverse_fill_ratio=result.adverse_fill_count / max(result.order_count, 1), venue=venue_tag, + fill_rate=result.fill_ratio, + slippage_bps=result.avg_slippage_bps, + price_improvement_bps=result.avg_price_improvement_bps, + fill_value_score=result.avg_fill_value_score, ) score = self._robust_score(results, params) return score, results @@ -1108,6 +1116,11 @@ class PolicyEvaluator: entropy_sum = 0.0 equity_start = state.account.equity cancel_count = 0 + # Fill quality accumulation (PRIMARY metrics) + fq_fill_value_sum = 0.0 + fq_price_improve_sum = 0.0 + fq_adverse_sum = 0.0 + fq_count = 0 for step in range(scenario.max_steps): # Plan with minimal overhead @@ -1145,6 +1158,14 @@ class PolicyEvaluator: cp_actions = tuple(cp.rollout_action(state, rng) for cp in scenario.counterparties) next_state = cwm.transition(state, (action, *cp_actions)) + # Accumulate fill quality from transition + if next_state.fill_quality: + fq = next_state.fill_quality + fq_fill_value_sum += fq.fill_value_score + fq_price_improve_sum += fq.price_improvement_bps + fq_adverse_sum += fq.post_fill_adverse_bps + fq_count += 1 + pnl = next_state.account.equity - equity_start pnl_bps = 10_000.0 * pnl / max(equity_start, 1.0) total_pnl_bps = pnl_bps @@ -1159,6 +1180,7 @@ class PolicyEvaluator: state = next_state steps = step + 1 if scenario.max_steps > 0 else 0 + fq_n = max(fq_count, 1) return EpisodeResult( scenario_id=scenario.scenario_id, policy_version=params.version, seed=rng_seed, steps=steps, pnl_bps=total_pnl_bps, realized_pnl=0.0, @@ -1172,6 +1194,9 @@ class PolicyEvaluator: policy_entropy_avg=entropy_sum / max(steps, 1), final_equity=state.account.equity, max_position_qty=0.0, diagnostics={"scenario_tags": scenario.tags}, + avg_fill_value_score=fq_fill_value_sum / fq_n, + avg_price_improvement_bps=fq_price_improve_sum / fq_n, + avg_post_fill_adverse_bps=fq_adverse_sum / fq_n, ) def _robust_score(self, results: list[EpisodeResult], params: FulfilmentPolicyParams) -> float: diff --git a/MALKHUT/malkhut/training/selector.py b/MALKHUT/malkhut/training/selector.py index ef49dbd..94b5f6e 100644 --- a/MALKHUT/malkhut/training/selector.py +++ b/MALKHUT/malkhut/training/selector.py @@ -147,8 +147,9 @@ class RegimeClassifier: @dataclass class RegimeStrategyScore: - """Performance score for a strategy in a specific regime. + """Performance score for a strategy in a specific regime + venue. + Fill quality metrics are the PRIMARY optimization target. Manifold fields (for Mode 2 recommendation): confidence: 0.0-1.0, how reliable is this score support_count: how many evaluations produced this score @@ -165,6 +166,11 @@ class RegimeStrategyScore: confidence: float = 1.0 support_count: int = 1 distance_to_nearest: float = 0.0 + # Fill quality metrics (PRIMARY) + avg_fill_rate: float = 0.0 + avg_slippage_bps: float = 0.0 + avg_price_improvement_bps: float = 0.0 + avg_fill_value_score: float = 0.0 class PerformanceMatrix: @@ -192,6 +198,10 @@ class PerformanceMatrix: drawdown_bps: float = 0.0, adverse_fill_ratio: float = 0.0, venue: str = "bingx", + fill_rate: float = 0.0, + slippage_bps: float = 0.0, + price_improvement_bps: float = 0.0, + fill_value_score: float = 0.0, ) -> None: """Record a strategy's performance in a regime on a specific venue.""" key = (regime, strategy_id, venue) @@ -204,12 +214,20 @@ class PerformanceMatrix: new_pnl = alpha * pnl_bps + (1 - alpha) * existing.avg_pnl_bps new_dd = alpha * drawdown_bps + (1 - alpha) * existing.avg_drawdown_bps new_adverse = alpha * adverse_fill_ratio + (1 - alpha) * existing.avg_adverse_fill_ratio + new_fill_rate = alpha * fill_rate + (1 - alpha) * existing.avg_fill_rate + new_slip = alpha * slippage_bps + (1 - alpha) * existing.avg_slippage_bps + new_improve = alpha * price_improvement_bps + (1 - alpha) * existing.avg_price_improvement_bps + new_fv = alpha * fill_value_score + (1 - alpha) * existing.avg_fill_value_score else: new_score = score new_episodes = 1 new_pnl = pnl_bps new_dd = drawdown_bps new_adverse = adverse_fill_ratio + new_fill_rate = fill_rate + new_slip = slippage_bps + new_improve = price_improvement_bps + new_fv = fill_value_score self._scores[key] = RegimeStrategyScore( strategy_id=strategy_id, @@ -223,6 +241,10 @@ class PerformanceMatrix: confidence=min(1.0, new_episodes / 10.0), support_count=new_episodes, distance_to_nearest=0.0, + avg_fill_rate=new_fill_rate, + avg_slippage_bps=new_slip, + avg_price_improvement_bps=new_improve, + avg_fill_value_score=new_fv, ) self._strategy_regime_history[strategy_id].append(regime)