Files
sentiment-engine/MALKHUT/malkhut/tests/test_e2e_integration.py
Codex d24d9bc6bd malkhut(wire): OrderType as three orthogonal dimensions — Fable's corrections
CRITICAL REFACTOR based on Fable's review (S9 roadmap item):

Before: flat enum conflating order types with TIF/instructions
  OrderType had MARKET, LIMIT, IOC, FOK, POST_ONLY, REDUCE_ONLY, etc.

After: three orthogonal dimensions (FIX-aligned):
  1. OrderType (Tag 40): what the order IS
     LIMIT, MARKET, STOP_MARKET, STOP_LIMIT, TRIGGER_MARKET, TRIGGER_LIMIT,
     TRAILING_STOP, OCO, TP_SL
  2. TimeInForce (Tag 59): how long it LIVES
     GTC, IOC, FOK, GTD
  3. Instructions (Tag 18): behavioral modifiers
     POST_ONLY, REDUCE_ONLY, HIDDEN, ICEBERG

Key corrections:
- POST_ONLY is an instruction on a LIMIT order, not a standalone type
- IOC/FOK are TimeInForce values, not order types
- BingX trailing_stop -> native TRAILING_STOP_MARKET (not TRIGGER_MARKET)
- FulfilmentAction.time_in_force: new field, default GTC

Exchange mappings restructured:
  EXCHANGE_ORDER_TYPE_MAP: OrderType -> exchange native 'type' param
  EXCHANGE_TIF_MAP: TimeInForce -> exchange native 'timeInForce' param
  EXCHANGE_INSTRUCTION_MAP: Instruction -> exchange encoding

21 files changed. 380+ tests pass. Backward compatible.
2026-07-14 14:46:44 +02:00

413 lines
21 KiB
Python

"""
End-to-end integration test — wire ALL subsystems together.
Simulates a full training cycle:
1. Create initial state + scenarios
2. Run training pipeline (CMA-ES → evaluate → promote)
3. Register promoted policy
4. Load into engine
5. Engine plans on live state
6. Risk gate validates
7. BingX adapter tracks order
8. Zinc SHM publishes book/account/fulfilment
9. ClickHouse persists decisions
10. Control plane sends HOT_RELOAD
11. Engine hot-reloads new policy
12. Replay verifier validates CWM against trajectory
13. Training logger records everything
"""
import json
import time
import pytest
from malkhut.state import (
AccountState, ExecutionIntent, FulfilmentPolicyParams, IntentKind,
MarketWorldState, Mode, OrderBookState, PositionState, PriceLevel,
Side, VenueRules,
)
from malkhut.actions import (
ActionKind, FulfilmentAction, PlannedPolicy, RiskDecision,
)
from malkhut.cwm.core import MinimalCryptoLOBCWM
from malkhut.cwm.replay_verify import (
ReplayVerifier, ReplayStep, TrajectoryRecorder,
)
from malkhut.planner.sm_mcts import DecoupledUCBPlanner
from malkhut.counterparties import default_counterparty_ecology
from malkhut.risk.gate import RiskGate
from malkhut.venue.bingx.adapter import BingXVenueAdapter, BingXConfig
from malkhut.ipc.zinc_plane import MalkhutZincPlane
from malkhut.ipc.control_plane import MalkhutControlPlane, ControlCommand, ControlPlaneFrame
from malkhut.storage.ch_store import MalkhutCHStore
from malkhut.execution.asex_integration import FulfilmentWorker, RiskWorker, RiskCheck
from malkhut.training.pipeline import TrainingPipeline, PipelineConfig, TrainingLogger
from malkhut.training.registry import PolicyRegistry, PolicyStage
from malkhut.training.cma_trainer import (
CMAParameterCodec, PolicyEvaluator, ScenarioFactory, SelfPlayPool,
PolicySnapshot,
)
# ── Helpers ──────────────────────────────────────────────────────────────────
def _venue(**kw):
d = dict(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)
d.update(kw)
return VenueRules(**d)
def _baseline(**kw):
d = dict(
version="baseline", 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.10, 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,
)
d.update(kw)
return FulfilmentPolicyParams(**d)
def _make_state(bid=50000.0, ask=50001.0, equity=10000.0):
return MarketWorldState(
ts_ns=1_000_000_000, mode=Mode.ENDOGENOUS_AGENT_SIM,
venue=_venue(),
book=OrderBookState(
ts_ns=1_000_000_000, symbol="BTCUSDT",
bids=(PriceLevel(bid, 1.0), PriceLevel(bid - 1.0, 2.0)),
asks=(PriceLevel(ask, 1.0), PriceLevel(ask + 1.0, 2.0)),
),
account=AccountState(
ts_ns=1_000_000_000, equity=equity, wallet_balance=equity,
available_balance=equity, margin_used=0.0, total_notional=0.0,
),
)
def _make_intent(symbol="BTCUSDT", urgency=0.5):
return ExecutionIntent(
intent_id=f"e2e_{int(time.time_ns())}", ts_ns=1_000_000_000,
symbol=symbol, kind=IntentKind.ENTER_LONG, target_qty=0.01,
max_notional=500.0, urgency=urgency, 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="e2e_test",
)
# ══════════════════════════════════════════════════════════════════════════════
# THE FULL END-TO-END TEST
# ══════════════════════════════════════════════════════════════════════════════
class TestEndToEndFullCycle:
"""
Wire ALL subsystems together and simulate a full training cycle:
training pipeline → registry → engine → plan → risk → venue →
zinc → clickhouse → control plane → hot reload → replay verify
"""
def test_full_e2e_cycle(self):
# ──────────────────────────────────────────────────────────────────
# 1. SETUP: All subsystems
# ──────────────────────────────────────────────────────────────────
zinc = MalkhutZincPlane(prefix="e2e_test")
control_plane = MalkhutControlPlane()
registry = PolicyRegistry()
engine = None
try:
# ──────────────────────────────────────────────────────────────
# 2. TRAINING: Run pipeline to produce a trained policy
# ──────────────────────────────────────────────────────────────
pipeline_config = PipelineConfig(
max_generations=1, max_evals_per_generation=7,
max_time_s=30, auto_promote=True,
)
pipeline = TrainingPipeline(
config=pipeline_config, registry=registry,
log_path="/dev/null",
)
pipeline_result = pipeline.run(
incumbent=_baseline(version="init"),
symbols=("BTCUSDT",),
)
# Verify training produced results
assert pipeline_result.generations_run >= 1
assert pipeline_result.total_evals > 0
assert len(pipeline_result.events) > 0
# ──────────────────────────────────────────────────────────────
# 3. REGISTRY: Check promoted policy exists
# ──────────────────────────────────────────────────────────────
active_params = registry.load_active()
assert active_params is not None, "No active policy after training"
trained_version = active_params.version
# ──────────────────────────────────────────────────────────────
# 4. ENGINE: Create and load trained policy
# ──────────────────────────────────────────────────────────────
from malkhut.engine import FulfilmentEngine
engine = FulfilmentEngine(
zinc=zinc, control_plane=control_plane,
registry=registry,
)
# Verify engine loaded the trained policy
current_params = engine.params_provider()
assert current_params.version == trained_version
# ──────────────────────────────────────────────────────────────
# 5. LIVE PLANNING: Engine plans on live state
# ──────────────────────────────────────────────────────────────
state = _make_state()
state_with_intent = MarketWorldState(
ts_ns=state.ts_ns, mode=state.mode, venue=state.venue,
book=state.book, account=state.account,
intent=_make_intent(),
)
planner = DecoupledUCBPlanner(
cwm=MinimalCryptoLOBCWM(),
counterparties=default_counterparty_ecology(),
rng_seed=42,
)
planned = planner.plan(
root_state=state_with_intent, params=current_params, budget_ms=20,
)
# Verify planner produced valid output
assert isinstance(planned, PlannedPolicy)
assert len(planned.actions) > 0
assert abs(sum(planned.probabilities) - 1.0) < 1e-6
assert planned.selected_action is not None
# ──────────────────────────────────────────────────────────────
# 6. RISK GATE: Validate the planned action
# ──────────────────────────────────────────────────────────────
risk_gate = RiskGate()
decision = risk_gate.validate(state_with_intent, planned, current_params)
assert isinstance(decision, RiskDecision)
# ──────────────────────────────────────────────────────────────
# 7. VENUE ADAPTER: Track the order
# ──────────────────────────────────────────────────────────────
adapter = BingXVenueAdapter()
adapter.execute(state_with_intent, decision)
# Verify adapter tracked the order (if approved)
if decision.approved and decision.action.kind != ActionKind.NOOP:
assert adapter.total_orders >= 1
working = adapter.get_working()
assert len(working) >= 1
# ──────────────────────────────────────────────────────────────
# 8. ZINC SHM: Publish book/account/fulfilment
# ──────────────────────────────────────────────────────────────
zinc.publish_book({
"ts_ns": state.ts_ns, "symbol": "BTCUSDT",
"bid": state.book.best_bid, "ask": state.book.best_ask,
})
book_data, book_seq = zinc.read_book()
assert book_data["symbol"] == "BTCUSDT"
assert book_seq >= 1
zinc.publish_account({
"ts_ns": state.ts_ns, "equity": state.account.equity,
})
acct_data, acct_seq = zinc.read_account()
assert acct_data["equity"] == 10000.0
zinc.publish_fulfilment({
"ts_ns": state.ts_ns, "action": str(planned.selected_action.kind.value),
"approved": decision.approved,
})
fulfil_data, fulfil_seq = zinc.read_fulfilment()
assert fulfil_seq >= 1
# ──────────────────────────────────────────────────────────────
# 9. CLICKHOUSE: Persist decision
# ──────────────────────────────────────────────────────────────
store = MalkhutCHStore()
store.ensure_tables()
store.store_fulfilment_decision(
ts_ns=state.ts_ns, exchange="bingx", symbol="BTCUSDT",
intent_id="e2e_test", state_hash="abc123",
selected_action=str(planned.selected_action.kind.value),
root_distribution=str(planned.probabilities),
risk_decision=f"{decision.approved}:{decision.reason}",
policy_version=current_params.version, latency_ms=5.0,
)
# Verify CH query
result = store.query("SELECT count() FROM fulfilment_decisions")
assert int(result.strip()) >= 1
# ──────────────────────────────────────────────────────────────
# 10. CONTROL PLANE: HOT_RELOAD_POLICY
# ──────────────────────────────────────────────────────────────
# Register a new policy and promote
new_params = _baseline(version="v2_reloaded")
registry.register_candidate(new_params, score=15.0)
registry.promote(new_params.version, PolicyStage.ACTIVE, "e2e reload")
# Send HOT_RELOAD via control plane
control_plane.publish_command(ControlPlaneFrame(
command=ControlCommand.HOT_RELOAD_POLICY.value,
ts_ns=time.time_ns(),
params={"policy_version": new_params.version},
source="e2e_test",
))
# Engine processes control plane
engine._process_control_plane()
# Verify engine reloaded
assert engine.params_provider().version == new_params.version
# ──────────────────────────────────────────────────────────────
# 11. REPLAY VERIFICATION: Validate CWM determinism
# ──────────────────────────────────────────────────────────────
cwm = MinimalCryptoLOBCWM()
recorder = TrajectoryRecorder(max_steps=10)
s = _make_state()
actions = [_noop(), _cross(Side.BUY, 0.1), _noop()]
for i, a in enumerate(actions):
after = cwm.transition(s, (a,))
recorder.record(i, s, (a,), after)
s = after
# Verify deterministic re-run
ok, mismatches = recorder.verify_deterministic(cwm)
assert ok, f"Determinism failed: {mismatches}"
# Verify against replay
verifier = ReplayVerifier()
replay = recorder.to_replay_steps()
result = verifier.verify(cwm, replay)
assert result.passed, f"Replay failed: {result.mismatches}"
# ──────────────────────────────────────────────────────────────
# 12. TRAINING LOGGER: Verify events were recorded
# ──────────────────────────────────────────────────────────────
events = pipeline.logger.get_events()
assert len(events) > 0
event_types = [e.event_type for e in events]
assert "run_start" in event_types
assert "generation" in event_types
assert "run_end" in event_types
# ──────────────────────────────────────────────────────────────
# 13. ASEX WORKERS: Verify state mutations through ASEx
# ──────────────────────────────────────────────────────────────
fw = engine.fulfilment_worker
assert fw.mutation_count >= 0
fw.reload_policy(new_params)
time.sleep(0.05)
assert fw.params is not None
# ──────────────────────────────────────────────────────────────
# 14. CLEANUP
# ──────────────────────────────────────────────────────────────
adapter.close()
engine.close()
finally:
zinc.close_all()
control_plane.close()
def test_multi_step_trajectory(self):
"""Run a multi-step trajectory through the full pipeline."""
cwm = MinimalCryptoLOBCWM()
recorder = TrajectoryRecorder(max_steps=20)
planner = DecoupledUCBPlanner(
cwm=cwm, counterparties=default_counterparty_ecology(), rng_seed=42,
)
params = _baseline()
risk_gate = RiskGate()
adapter = BingXVenueAdapter()
zinc = MalkhutZincPlane(prefix="e2e_traj")
try:
s = _make_state()
total_pnl = 0.0
for step in range(10):
# Add intent
intent = _make_intent(urgency=0.5)
s_with_intent = MarketWorldState(
ts_ns=s.ts_ns, mode=s.mode, venue=s.venue,
book=s.book, account=s.account, intent=intent,
)
# Plan
planned = planner.plan(root_state=s_with_intent, params=params, budget_ms=15)
# Risk gate
decision = risk_gate.validate(s_with_intent, planned, params)
# Venue adapter
adapter.execute(s_with_intent, decision)
# Zinc
zinc.publish_book({"ts_ns": s.ts_ns, "symbol": "BTCUSDT"})
# CWM transition
cps = default_counterparty_ecology()
import random
rng = random.Random(42 + step)
cp_actions = tuple(cp.rollout_action(s, rng) for cp in cps)
next_state = cwm.transition(s, (planned.selected_action, *cp_actions))
# Record trajectory
recorder.record(step, s, (planned.selected_action, *cp_actions), next_state)
# Track PnL
pnl = next_state.account.equity - s.account.equity
total_pnl += pnl
s = next_state
# Verify trajectory
ok, mismatches = recorder.verify_deterministic(cwm)
assert ok, f"Determinism failed at step: {mismatches}"
# Verify replay
verifier = ReplayVerifier()
replay = recorder.to_replay_steps()
result = verifier.verify(cwm, replay)
assert result.passed
# Verify final state is valid
assert s.account.equity > 0
assert s.ts_ns > 1_000_000_000
finally:
adapter.close()
zinc.close_all()
# ── Helpers (local) ──────────────────────────────────────────────────────────
def _noop():
return FulfilmentAction(ActionKind.NOOP, None, None, 0, 0.0, 0)
def _cross(side, frac=0.1):
from malkhut.actions import OrderType
return FulfilmentAction(ActionKind.CROSS_SPREAD, side, OrderType.LIMIT, 0, frac, 50)