diff --git a/MALKHUT/malkhut/tests/test_scoring_modes.py b/MALKHUT/malkhut/tests/test_scoring_modes.py new file mode 100644 index 0000000..e61ca37 --- /dev/null +++ b/MALKHUT/malkhut/tests/test_scoring_modes.py @@ -0,0 +1,116 @@ +""" +Tests for scoring modes: fast scalar + advantage estimation. +""" +import pytest +from malkhut.training.cma_trainer import ( + ScenarioFactory, PolicyEvaluator, EpisodeResult +) +from malkhut.cwm.core import MinimalCryptoLOBCWM +from malkhut.state import FulfilmentPolicyParams + + +def _make_params(max_sims=16): + return FulfilmentPolicyParams( + version='test', ucb_c=1.414, max_sims=max_sims, max_depth=2, + rollout_depth=2, root_temperature=0.5, min_root_entropy=0.25, + quote_offsets_ticks=(0, 1), quote_size_fractions=(0.25, 0.50), + passive_ttl_ms=200, aggressive_ttl_ms=50, + maker_edge_min_bps=0.5, cross_spread_edge_min_bps=5.0, + adverse_toxicity_cancel_threshold=0.5, queue_churn_cancel_threshold=0.5, + mae_tail_cut_bps=50.0, mfe_giveback_cut_fraction=0.5, + max_time_in_loss_s=300.0, failed_recovery_cut_count=3, + recovery_velocity_min_bps_per_s=0.0, + max_symbol_notional_fraction=0.20, max_single_order_notional_fraction=0.05, + reduce_when_global_up_fraction=0.30, session_profit_lock_fraction=0.02, + w_expected_pnl=1.0, w_fill_probability=0.5, w_adverse_selection=2.0, + w_queue_priority=0.5, w_inventory_risk=1.5, w_tail_loss=5.0, + w_fee_quality=0.5, w_time_decay=0.3, w_policy_entropy=0.5, + robust_tail_weight=2.0, toxic_counterparty_weight=3.0, + low_liquidity_weight=2.0, latency_stress_weight=1.0, + ) + + +class TestFastScoring: + def test_fast_mode_default(self): + evaluator = PolicyEvaluator(cwm_factory=MinimalCryptoLOBCWM) + assert evaluator._scoring_mode == "fast" + + def test_advantage_mode_init(self): + evaluator = PolicyEvaluator(cwm_factory=MinimalCryptoLOBCWM, scoring_mode="advantage") + assert evaluator._scoring_mode == "advantage" + + def test_fast_score_with_fills(self): + evaluator = PolicyEvaluator(cwm_factory=MinimalCryptoLOBCWM) + factory = ScenarioFactory() + suite = factory.build_suite(symbols=('BTCUSDT',), steps_per_scenario=5) + params = _make_params() + score, results = evaluator.evaluate_candidate( + params=params, scenarios=suite[:5], rng_seed=42, workers=0) + assert isinstance(score, float) + assert score > -1000 + + def test_fast_score_empty_results(self): + evaluator = PolicyEvaluator(cwm_factory=MinimalCryptoLOBCWM) + score = evaluator._robust_score([], _make_params()) + assert score == -1000.0 + + def test_fast_score_bounded(self): + evaluator = PolicyEvaluator(cwm_factory=MinimalCryptoLOBCWM) + factory = ScenarioFactory() + suite = factory.build_suite(symbols=('BTCUSDT', 'ETHUSDT'), steps_per_scenario=5) + params = _make_params() + scores = [] + for i in range(20): + score, _ = evaluator.evaluate_candidate( + params=params, scenarios=suite[:5], rng_seed=i, workers=0) + scores.append(score) + for s in scores: + assert isinstance(s, float) + assert s > -10000 # no NaN or inf + + +class TestAdvantageScoring: + def test_advantage_mode(self): + evaluator = PolicyEvaluator( + cwm_factory=MinimalCryptoLOBCWM, scoring_mode="advantage") + factory = ScenarioFactory() + suite = factory.build_suite(symbols=('BTCUSDT',), steps_per_scenario=5) + params = _make_params() + score, results = evaluator.evaluate_candidate( + params=params, scenarios=suite[:3], rng_seed=42, workers=0) + assert isinstance(score, float) + assert -10.0 <= score <= 10.0 + + def test_advantage_baseline_evolution(self): + evaluator = PolicyEvaluator( + cwm_factory=MinimalCryptoLOBCWM, scoring_mode="advantage") + factory = ScenarioFactory() + suite = factory.build_suite(symbols=('BTCUSDT',), steps_per_scenario=3) + params = _make_params() + baselines = [] + for i in range(10): + evaluator.evaluate_candidate( + params=params, scenarios=suite[:2], rng_seed=i, workers=0) + baselines.append(evaluator._adv_baseline) + # Baseline should converge (not oscillate wildly) + last_5 = baselines[-5:] + assert max(last_5) - min(last_5) < abs(baselines[0]) + 1.0 + + +class TestScoringModeIntegration: + def test_evaluate_candidate_with_mode(self): + factory = ScenarioFactory() + suite = factory.build_suite(symbols=('BTCUSDT',), steps_per_scenario=3) + params = _make_params() + + # Fast mode + ev_fast = PolicyEvaluator(cwm_factory=MinimalCryptoLOBCWM, scoring_mode="fast") + score_f, _ = ev_fast.evaluate_candidate(params=params, scenarios=suite[:2], rng_seed=0, workers=0) + + # Advantage mode + ev_adv = PolicyEvaluator(cwm_factory=MinimalCryptoLOBCWM, scoring_mode="advantage") + score_a, _ = ev_adv.evaluate_candidate(params=params, scenarios=suite[:2], rng_seed=0, workers=0) + + # Both should produce valid scores + assert isinstance(score_f, float) + assert isinstance(score_a, float) diff --git a/MALKHUT/malkhut/training/advantage_scorer.py b/MALKHUT/malkhut/training/advantage_scorer.py new file mode 100644 index 0000000..c14dca2 --- /dev/null +++ b/MALKHUT/malkhut/training/advantage_scorer.py @@ -0,0 +1,110 @@ +""" +Advantage-based scoring for MALKHUT CMA-ES training. + +Replaces raw reward with advantage estimation: + advantage(action) = actual_return - baseline_return + +Where baseline = running mean of recent returns. This lets the system learn: + - WHEN to trade (positive advantage → conditions were favorable) + - WHEN NOT to trade (negative advantage → conditions were unfavorable) + - WHY orders don't fill (reason tracking → next action) + +Knobs: + baseline_decay: how fast baseline adapts (default 0.995) + advantage_clip: clip extreme values (default ±10.0) + fill_weight: reward weight for fills (default 1.0) + counterfactual_weight: weight for missed-opportunity penalty (default 0.5) +""" +from __future__ import annotations + +from dataclasses import dataclass, field +from typing import Dict, List, Optional, Sequence +from malkhut.training.cma_trainer import EpisodeResult + + +@dataclass +class AdvantageScorer: + """Tracks baseline and computes advantage scores for CMA-ES.""" + baseline_decay: float = 0.995 + advantage_clip: float = 10.0 + fill_weight: float = 1.0 + counterfactual_weight: float = 0.5 + reason_weight: float = 2.0 + + # Running state + _baseline: float = field(default=0.0, init=False) + _n_seen: int = field(default=0, init=False) + _total_score: float = field(default=0.0, init=False) + + def reset(self) -> None: + self._baseline = 0.0 + self._n_seen = 0 + self._total_score = 0.0 + + def score(self, results: Sequence[EpisodeResult]) -> float: + """Compute advantage-based score for a batch of episodes. + + Instead of penalizing no-ops or rewarding raw fills, this scores + based on: "was this batch better or worse than average?" + """ + if not results: + return -1000.0 + + # Raw performance metrics + total_pnl = sum(r.pnl_bps for r in results) + n_fills = sum(r.fill_count for r in results) + n_noops = sum(r.noop_count for r in results) + n_orders = sum(r.order_count for r in results) + n_episodes = len(results) + + # Fill quality (adverse selection) + total_adverse = sum(r.adverse_fill_count for r in results) + fill_ratio = n_fills / max(n_orders, 1) + adverse_ratio = total_adverse / max(n_fills, 1) + + # Drawdown penalty + avg_dd = sum(r.max_drawdown_bps for r in results) / n_episodes + + # Entropy (diversity of actions) + avg_entropy = sum(r.policy_entropy_avg for r in results) / n_episodes + + # Raw score (weighted combination) + raw_score = ( + total_pnl * 10.0 # PnL is king + + n_fills * self.fill_weight * 5.0 # reward fills + - adverse_ratio * 20.0 # penalize adverse selection + - avg_dd * 2.0 # penalize drawdown + + avg_entropy * 0.1 # reward diversity + ) + + # Reason tracking: learning from unfilled orders + unfilled = n_orders - n_fills + if unfilled > 0 and n_orders > 0: + # Unfilled orders that were NOT noops = the planner tried but failed + # This is information: we can learn from it + unfilled_ratio = unfilled / n_orders + # Penalize UNNECESSARY unfilled orders (placed but too far from market) + # But NOT too heavily — some unfilled is normal + raw_score -= unfilled_ratio * self.reason_weight * 3.0 + + # Advantage = raw - baseline + advantage = raw_score - self._baseline + + # Clip extreme advantages + advantage = max(-self.advantage_clip, min(self.advantage_clip, advantage)) + + # Update baseline (exponential moving average) + if self._n_seen == 0: + self._baseline = raw_score + else: + self._baseline = self.baseline_decay * self._baseline + (1 - self.baseline_decay) * raw_score + self._n_seen += 1 + self._total_score += raw_score + + return advantage + + def get_baseline(self) -> float: + return self._baseline + + def get_mean_score(self) -> float: + return self._total_score / max(self._n_seen, 1) diff --git a/MALKHUT/malkhut/training/cma_trainer.py b/MALKHUT/malkhut/training/cma_trainer.py index e176cc0..a67098f 100644 --- a/MALKHUT/malkhut/training/cma_trainer.py +++ b/MALKHUT/malkhut/training/cma_trainer.py @@ -951,9 +951,13 @@ class PolicyEvaluator: self, cwm_factory: Callable[[], CodeWorldModel], counterparties: Optional[Tuple[CounterpartyPolicy, ...]] = None, + scoring_mode: str = "fast", ) -> None: self.cwm_factory = cwm_factory self.counterparties = counterparties or default_counterparty_ecology() + self._scoring_mode = scoring_mode + self._adv_baseline = 0.0 + self._adv_n_seen = 0 def evaluate_candidate( self, @@ -1130,43 +1134,111 @@ class PolicyEvaluator: ) def _robust_score(self, results: list[EpisodeResult], params: FulfilmentPolicyParams) -> float: + """Score execution quality. + + Modes: + "fast" (default): simple scalar for CMA loop — rewards good fills, + tolerates no-fills (valid advisory), penalizes extremes. + "advantage": full advantage estimation for offline analysis. + """ if not results: - return -float("inf") + return -1000.0 - pnl = [r.pnl_bps for r in results] - pnl_sorted = sorted(pnl) - tail_idx = max(0, int(TAIL_QUANTILE * (len(pnl_sorted) - 1))) - p05 = pnl_sorted[tail_idx] - mean = sum(pnl) / len(pnl) + if getattr(self, '_scoring_mode', 'fast') == 'advantage': + return self._advantage_score(results) - adverse = sum(r.adverse_fill_count for r in results) / max(sum(r.order_count for r in results), 1) - slippage = sum(r.avg_slippage_bps for r in results) / len(results) - liq = sum(r.liquidation_near_miss_count for r in results) - dd = sum(r.max_drawdown_bps for r in results) / len(results) - entropy = sum(r.policy_entropy_avg for r in results) / len(results) + # === FAST SCALAR MODE === + # Reward: execution quality (good fills, fast fills, low adverse selection) + # Tolerate: no-fills (valid advisory recommendation) + # Penalize: extreme fill rates, adverse selection, drawdown - score = 0.0 - score += mean * 10.0 # HEAVY PnL weight - score += params.robust_tail_weight * p05 * 5.0 - score -= params.toxic_counterparty_weight * adverse * 100.0 - score -= slippage - score -= 10.0 * liq - score -= 2.0 * dd # penalize drawdown - score += params.w_policy_entropy * entropy * 0.1 # reduced entropy weight + n_episodes = len(results) + n_fills = sum(r.fill_count for r in results) + n_orders = sum(r.order_count for r in results) + n_noops = sum(r.noop_count for r in results) - # NOOP penalty: penalize strategies that don't trade - noop_ratios = [r.noop_count / max(r.steps, 1) for r in results] - avg_noop_ratio = sum(noop_ratios) / len(noop_ratios) if noop_ratios else 0.0 - score -= avg_noop_ratio * 50.0 # heavy penalty for not trading + # --- Execution quality: PnL when fills happen --- + fill_pnls = [r.pnl_bps for r in results if r.fill_count > 0] + if fill_pnls: + mean_fill_pnl = sum(fill_pnls) / len(fill_pnls) + else: + mean_fill_pnl = 0.0 - # Fill reward: reward strategies that actually get fills - fill_ratios = [r.fill_count / max(r.order_count, 1) for r in results] - avg_fill_ratio = sum(fill_ratios) / len(fill_ratios) if fill_ratios else 0.0 - score += avg_fill_ratio * 20.0 # reward fills + # --- Fill rate: reward moderate, penalize extremes --- + fill_rate = n_fills / max(n_orders, 1) + # Sweet spot: 5-15% fill rate → bonus + # Too low (<3%): not enough trading → small penalty + # Too high (>30%): getting picked off → heavy penalty + if fill_rate < 0.03: + fill_bonus = -2.0 * (0.03 - fill_rate) / 0.03 # penalty for too few fills + elif fill_rate > 0.30: + fill_bonus = -5.0 * (fill_rate - 0.30) / 0.70 # penalty for too many fills + else: + fill_bonus = 2.0 * (fill_rate - 0.03) / 0.12 # bonus in sweet spot (0-2 points) + + # --- Adverse selection --- + total_adverse = sum(r.adverse_fill_count for r in results) + adverse_ratio = total_adverse / max(n_fills, 1) + adverse_penalty = -3.0 * adverse_ratio + + # --- Drawdown --- + avg_dd = sum(r.max_drawdown_bps for r in results) / n_episodes + dd_penalty = -0.5 * avg_dd + + # --- Reason tracking: learn from unfilled orders --- + unfilled = n_orders - n_fills + noop_ratio = n_noops / max(n_episodes * 10, 1) # normalize by max steps + unfilled_ratio = unfilled / max(n_orders, 1) + # Light penalty for too many noops (but NOT heavy like before) + noop_penalty = -0.5 * noop_ratio + + # --- Total score --- + score = ( + mean_fill_pnl * 2.0 # execution quality when fills happen + + fill_bonus # reward moderate fill rate + + adverse_penalty # penalize adverse selection + + dd_penalty # penalize drawdown + + noop_penalty # light noop penalty + ) return score - return score + def _advantage_score(self, results: list[EpisodeResult]) -> float: + """Full advantage estimation for offline analysis. + + advantage = raw_performance - baseline_performance + baseline = exponential moving average of recent raw scores. + """ + n_episodes = len(results) + n_fills = sum(r.fill_count for r in results) + n_orders = sum(r.order_count for r in results) + total_adverse = sum(r.adverse_fill_count for r in results) + avg_dd = sum(r.max_drawdown_bps for r in results) / n_episodes + avg_entropy = sum(r.policy_entropy_avg for r in results) / n_episodes + fill_pnls = [r.pnl_bps for r in results if r.fill_count > 0] + mean_fill_pnl = sum(fill_pnls) / len(fill_pnls) if fill_pnls else 0.0 + + # Raw performance + raw = ( + mean_fill_pnl * 10.0 + + n_fills * 5.0 + - (total_adverse / max(n_fills, 1)) * 20.0 + - avg_dd * 2.0 + + avg_entropy * 0.1 + ) + + # Update baseline + if not hasattr(self, '_adv_baseline'): + self._adv_baseline = 0.0 + if self._adv_baseline == 0.0 and self._adv_n_seen == 0: + self._adv_baseline = raw + else: + self._adv_baseline = 0.995 * self._adv_baseline + 0.005 * raw + self._adv_n_seen += 1 + + # Advantage = raw - baseline, clipped + advantage = raw - self._adv_baseline + return max(-10.0, min(10.0, advantage)) @staticmethod def performance_vector(results: list[EpisodeResult]) -> Tuple[float, ...]: