Files
sentiment-engine/MALKHUT/malkhut/tests/test_harness.py
Codex 4c239f7774 malkhut(tests): 1140 test functions across 46 test files
CWM (103): core mechanics, exhaustive edge cases, numba, exchange mechanics
Replay (118): exhaustive verification, microstructure, trajectory
Training (190): asset classification, phase0 extensive, pipeline, exhaustive
DSL (102): v2 syntax, expanded, new features
ASEx (33): validate-before-mutate, single-writer
Planner (48): MCTS, alternatives, hooks
Counterparties (19): 9 adversarial agent policies
Clock (30): event-driven reactor
BingX (28): venue adapter
IPC (8): Zinc SHM
Storage (9): ClickHouse
Risk (4): hard invariants
State (17): frozen dataclass invariants
Integration: E2E, concurrency, sync/async seams, hypothesis, fuzz, adversarial
2026-07-11 10:46:12 +02:00

367 lines
18 KiB
Python

"""
Harness tests — verify the system CAN learn and CAN register improvement.
These tests prove the system is CAPABLE of learning:
1. Different parameters produce different actions
2. Different actions produce different PnL
3. CMA-ES can find better parameters
4. Genetic operators produce meaningful diversity
5. Score improves over generations
"""
import random
import pytest
from malkhut.state import (
AccountState, ExecutionIntent, FulfilmentPolicyParams, IntentKind,
MarketWorldState, Mode, OrderBookState, PriceLevel, Side, VenueRules,
)
from malkhut.actions import ActionKind, FulfilmentAction, OrderType
from malkhut.cwm.core import MinimalCryptoLOBCWM
from malkhut.counterparties import default_counterparty_ecology
from malkhut.planner.sm_mcts import DecoupledUCBPlanner
def _venue():
return VenueRules(exchange="bingx", symbol="BTCUSDT", tick_size=0.1, lot_size=0.001,
min_qty=0.001, min_notional=5.0, maker_fee_bps=-0.2, taker_fee_bps=0.5,
post_only_supported=True, reduce_only_supported=True,
max_orders_per_second=100, max_cancels_per_minute=120)
def _state():
return MarketWorldState(
ts_ns=1, mode=Mode.REPLAY_NO_IMPACT, venue=_venue(),
book=OrderBookState(ts_ns=1, symbol="BTCUSDT",
bids=(PriceLevel(50000.0, 1.0),),
asks=(PriceLevel(50001.0, 1.0),)),
account=AccountState(ts_ns=1, equity=10000.0, wallet_balance=10000.0,
available_balance=10000.0, margin_used=0.0, total_notional=0.0),
)
def _intent():
return ExecutionIntent(
intent_id="test", ts_ns=1, symbol="BTCUSDT",
kind=IntentKind.ENTER_LONG, target_qty=0.01, max_notional=500.0,
urgency=0.5, alpha_horizon_s=60.0, alpha_bps=2.0,
max_slippage_bps=5.0, prefer_maker=True, reduce_only=False,
ttl_s=300.0, reason="test",
)
def _params(**kw):
d = dict(
version="test", ucb_c=1.414, max_sims=64, max_depth=2, rollout_depth=2,
root_temperature=0.5, min_root_entropy=0.25, quote_offsets_ticks=(0, 1, 2),
quote_size_fractions=(0.1, 0.25, 0.5), 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,
)
d.update(kw)
return FulfilmentPolicyParams(**d)
# ══════════════════════════════════════════════════════════════════════════════
# HARNESS 1: Different parameters produce different actions
# ══════════════════════════════════════════════════════════════════════════════
class TestParameterSensitivity:
def test_ucb_c_affects_exploration(self):
"""Different UCB_c should produce different action distributions."""
cwm = MinimalCryptoLOBCWM()
s = _state()
intent = _intent()
s_with_intent = MarketWorldState(
ts_ns=1, mode=Mode.REPLAY_NO_IMPACT, venue=_venue(),
book=OrderBookState(ts_ns=1, symbol="BTCUSDT",
bids=(PriceLevel(50000.0, 1.0),),
asks=(PriceLevel(50001.0, 1.0),)),
account=AccountState(ts_ns=1, equity=10000.0, wallet_balance=10000.0,
available_balance=10000.0, margin_used=0.0, total_notional=0.0),
intent=intent,
)
actions_low = []
actions_high = []
for seed in range(30): # more samples for reliability
p_low = DecoupledUCBPlanner(cwm=cwm, counterparties=default_counterparty_ecology(), rng_seed=seed)
r_low = p_low.plan(s_with_intent, _params(ucb_c=0.2), budget_ms=10)
actions_low.append(r_low.selected_action.kind)
p_high = DecoupledUCBPlanner(cwm=cwm, counterparties=default_counterparty_ecology(), rng_seed=seed)
r_high = p_high.plan(s_with_intent, _params(ucb_c=3.0), budget_ms=10)
actions_high.append(r_high.selected_action.kind)
# Different exploration should produce different distributions
low_types = set(actions_low)
high_types = set(actions_high)
# Either different action sets OR different distribution within same set
assert low_types != high_types or len(low_types) > 1
def test_temperature_affects_distribution(self):
"""Different temperatures should produce different probability distributions."""
cwm = MinimalCryptoLOBCWM()
s = _state()
intent = _intent()
s_with_intent = MarketWorldState(
ts_ns=1, mode=Mode.REPLAY_NO_IMPACT, venue=_venue(),
book=OrderBookState(ts_ns=1, symbol="BTCUSDT",
bids=(PriceLevel(50000.0, 1.0),),
asks=(PriceLevel(50001.0, 1.0),)),
account=AccountState(ts_ns=1, equity=10000.0, wallet_balance=10000.0,
available_balance=10000.0, margin_used=0.0, total_notional=0.0),
intent=intent,
)
probs_by_temp = {}
for temp in [0.1, 0.5, 1.0, 2.0]:
p = DecoupledUCBPlanner(cwm=cwm, counterparties=default_counterparty_ecology(), rng_seed=42)
r = p.plan(s_with_intent, _params(root_temperature=temp), budget_ms=10)
probs_by_temp[temp] = tuple(r.probabilities)
# Different temperatures should produce different distributions
unique_dists = set(probs_by_temp.values())
assert len(unique_dists) > 1
# ══════════════════════════════════════════════════════════════════════════════
# HARNESS 2: Different actions produce different PnL
# ══════════════════════════════════════════════════════════════════════════════
class TestActionPnLDifferentiation:
def test_cross_vs_noop_different_equity(self):
"""CROSS and NOOP should produce different equity."""
cwm = MinimalCryptoLOBCWM()
s = _state()
r_cross = cwm.transition(s, (_cross(Side.BUY, 0.1),))
r_noop = cwm.transition(s, (_noop(),))
assert r_cross.account.equity != r_noop.account.equity
def test_cross_vs_place_different_equity(self):
"""CROSS and PLACE should produce different equity."""
cwm = MinimalCryptoLOBCWM()
s = _state()
r_cross = cwm.transition(s, (_cross(Side.BUY, 0.1),))
r_place = cwm.transition(s, (_place(Side.BUY, offset=0, frac=0.1),))
assert r_cross.account.equity != r_place.account.equity
def test_different_cross_sizes_different_equity(self):
"""Different cross sizes should produce different equity."""
cwm = MinimalCryptoLOBCWM()
s = _state()
r_small = cwm.transition(s, (_cross(Side.BUY, 0.01),))
r_large = cwm.transition(s, (_cross(Side.BUY, 0.1),))
assert r_small.account.equity != r_large.account.equity
# ══════════════════════════════════════════════════════════════════════════════
# HARNESS 3: CMA-ES can find better parameters
# ══════════════════════════════════════════════════════════════════════════════
class TestCMAESLearning:
def test_cma_es_finds_better_params(self):
"""CMA-ES should find parameters that produce different (hopefully better) scores."""
from malkhut.training.cma_trainer import CMAESTrainer, CMAParameterCodec, PolicyEvaluator, ScenarioFactory
from malkhut.counterparties import default_counterparty_ecology
codec = CMAParameterCodec()
evaluator = PolicyEvaluator(
cwm_factory=lambda: MinimalCryptoLOBCWM(),
counterparties=default_counterparty_ecology(),
)
pool = SelfPlayPool(max_size=5)
trainer = CMAESTrainer(codec=codec, evaluator=evaluator, pool=pool)
factory = ScenarioFactory()
scenarios = factory.build_suite(symbols=("BTCUSDT",), steps_per_scenario=10)
# Run CMA-ES for a few evaluations
import cma
x0 = codec.initial_vector(_baseline())
lows, highs = codec.bounds()
es = cma.CMAEvolutionStrategy(x0, 0.30, {
"bounds": [lows, highs], "popsize": 5, "seed": 42, "verbose": -9,
})
scores = []
for _ in range(3):
xs = es.ask()
for x in xs:
candidate = codec.decode(x, version="test")
score, _ = evaluator.evaluate_candidate(
params=candidate, scenarios=scenarios, rng_seed=42,
)
scores.append(score)
es.tell(xs, [-s for s in scores[-5:]])
# Scores should vary (not all identical)
assert len(set(scores)) > 1, "All scores identical — system can't learn"
# ══════════════════════════════════════════════════════════════════════════════
# HARNESS 4: Genetic operators produce meaningful diversity
# ══════════════════════════════════════════════════════════════════════════════
class TestGeneticDiversity:
def test_crossover_produces_different_children(self):
"""Crossover of two parents should produce different offspring."""
from malkhut.training.generator import GeneticOperators, StrategyGenome
from malkhut.training.cma_trainer import CMAParameterCodec
codec = CMAParameterCodec()
ops = GeneticOperators(codec=codec)
p1 = StrategyGenome(strategy_type=StrategyType.SM_MCTS, params=_baseline(version="p1"))
p2 = StrategyGenome(strategy_type=StrategyType.UCB1, params=_baseline(version="p2"))
children = []
for i in range(10):
child = ops.crossover(p1, p2, random.Random(i))
children.append(child)
# Children should have different params
unique_versions = set(c.params.version for c in children)
assert len(unique_versions) > 1
def test_mutation_produces_different_children(self):
"""Mutation should produce different offspring."""
from malkhut.training.generator import GeneticOperators, StrategyGenome
from malkhut.training.cma_trainer import CMAParameterCodec
codec = CMAParameterCodec()
ops = GeneticOperators(codec=codec, mutation_rate=0.5) # high mutation
parent = StrategyGenome(strategy_type=StrategyType.SM_MCTS, params=_baseline())
children = []
for i in range(10):
child = ops.mutate(parent, random.Random(i))
children.append(child)
# Children should have different params
unique_versions = set(c.params.version for c in children)
assert len(unique_versions) > 1
# ══════════════════════════════════════════════════════════════════════════════
# HARNESS 5: Score improves over generations
# ══════════════════════════════════════════════════════════════════════════════
class TestScoreImprovement:
def test_score_varies_across_generations(self):
"""Scores should vary across generations (not all identical)."""
from malkhut.training.cma_trainer import CMAESTrainer, CMAParameterCodec, PolicyEvaluator, ScenarioFactory
from malkhut.counterparties import default_counterparty_ecology
codec = CMAParameterCodec()
evaluator = PolicyEvaluator(
cwm_factory=lambda: MinimalCryptoLOBCWM(),
counterparties=default_counterparty_ecology(),
)
pool = SelfPlayPool(max_size=5)
trainer = CMAESTrainer(codec=codec, evaluator=evaluator, pool=pool)
factory = ScenarioFactory()
scenarios = factory.build_suite(symbols=("BTCUSDT",), steps_per_scenario=5)
# Run for a few generations
import cma
x0 = codec.initial_vector(_baseline())
lows, highs = codec.bounds()
es = cma.CMAEvolutionStrategy(x0, 0.30, {
"bounds": [lows, highs], "popsize": 5, "seed": 42, "verbose": -9,
})
all_scores = []
for gen in range(3):
xs = es.ask()
gen_scores = []
for x in xs:
candidate = codec.decode(x, version=f"gen{gen}")
score, _ = evaluator.evaluate_candidate(
params=candidate, scenarios=scenarios, rng_seed=42 + gen,
)
gen_scores.append(score)
all_scores.append(max(gen_scores))
es.tell(xs, [-s for s in gen_scores])
# Scores should vary (not all identical)
assert len(set(all_scores)) > 1, f"All generation scores identical: {all_scores}"
# ══════════════════════════════════════════════════════════════════════════════
# HARNESS 6: End-to-end training produces improvement
# ══════════════════════════════════════════════════════════════════════════════
class TestEndToEndLearning:
def test_training_pipeline_produces_improvement(self):
"""Training pipeline should produce improvement over baseline."""
from malkhut.training.pipeline import TrainingPipeline, PipelineConfig
from malkhut.training.registry import PolicyRegistry
from malkhut.storage.ch_store import MalkhutCHStore
store = MalkhutCHStore()
store.ensure_tables()
registry = PolicyRegistry(store=store)
cfg = PipelineConfig(max_generations=3, max_evals_per_generation=5, max_time_s=30)
pipeline = TrainingPipeline(config=cfg, registry=registry, log_path="/dev/null")
result = pipeline.run(incumbent=_baseline(), symbols=("BTCUSDT",))
# Should have run some generations
assert result.generations_run >= 1
# Should have some events
assert len(result.events) > 0
# Best score should be a valid number
assert isinstance(result.best_score, float)
# ══════════════════════════════════════════════════════════════════════════════
# HELPERS
# ══════════════════════════════════════════════════════════════════════════════
from malkhut.training.generator import StrategyType, SelfPlayPool
def _baseline(**kw):
d = dict(
version="baseline", ucb_c=1.414, max_sims=256, max_depth=3,
rollout_depth=3, root_temperature=0.5, min_root_entropy=0.25,
quote_offsets_ticks=(0, 1, 2), quote_size_fractions=(0.1, 0.25, 0.5),
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,
)
d.update(kw)
return FulfilmentPolicyParams(**d)
def _noop():
return FulfilmentAction(ActionKind.NOOP, None, None, 0, 0.0, 0)
def _cross(side, frac):
return FulfilmentAction(ActionKind.CROSS_SPREAD, side, OrderType.IOC, 0, frac, 50)
def _place(side, offset=0, frac=0.1):
return FulfilmentAction(ActionKind.PLACE, side, OrderType.LIMIT, offset, frac, 200)