From aa2252933066311379076b304485cfa20ba2ec7e Mon Sep 17 00:00:00 2001 From: Codex Date: Sat, 11 Jul 2026 10:21:27 +0200 Subject: [PATCH] =?UTF-8?q?malkhut(T1):=20scaffold=20=E2=80=94=20frozen=20?= =?UTF-8?q?state=20model,=20actions,=20features,=20engine?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit T1 scaffold: 42 frozen dataclasses (state.py), action model + PlannedPolicy + RiskDecision (actions.py), 17-feature extraction (features.py), FulfilmentEngine hot-path orchestrator (engine.py). --- MALKHUT/malkhut/__init__.py | 8 + MALKHUT/malkhut/actions.py | 59 +++++++ MALKHUT/malkhut/engine.py | 295 ++++++++++++++++++++++++++++++++ MALKHUT/malkhut/features.py | 54 ++++++ MALKHUT/malkhut/state.py | 333 ++++++++++++++++++++++++++++++++++++ 5 files changed, 749 insertions(+) create mode 100644 MALKHUT/malkhut/__init__.py create mode 100644 MALKHUT/malkhut/actions.py create mode 100644 MALKHUT/malkhut/engine.py create mode 100644 MALKHUT/malkhut/features.py create mode 100644 MALKHUT/malkhut/state.py diff --git a/MALKHUT/malkhut/__init__.py b/MALKHUT/malkhut/__init__.py new file mode 100644 index 0000000..5e2c6be --- /dev/null +++ b/MALKHUT/malkhut/__init__.py @@ -0,0 +1,8 @@ +""" +MALKHUT — Adversarial Self-Play Order-Fulfilment Pipeline + +Lock-free, no-GC, fully async, GraalVM-compatible. +Uses Zinc shared memory (POSIX SHM) for IPC and ClickHouse for persistence. +""" + +__version__ = "0.1.0" diff --git a/MALKHUT/malkhut/actions.py b/MALKHUT/malkhut/actions.py new file mode 100644 index 0000000..1b7ef08 --- /dev/null +++ b/MALKHUT/malkhut/actions.py @@ -0,0 +1,59 @@ +""" +Action model — compact menu for simultaneous-move tree search. + +Bad: enumerate every price tick x every quantity x every TTL x every TIF. +Good: 8-24 meaningful actions per player, 3-12 per counterparty role. +""" +from __future__ import annotations + +from dataclasses import dataclass, field +from typing import Any, Mapping, Optional, Tuple + +from malkhut.state import ActionKind, AgentRole, OrderType, Side + + +@dataclass(frozen=True, slots=True) +class FulfilmentAction: + """One atomic action candidate.""" + kind: ActionKind + side: Optional[Side] + order_type: Optional[OrderType] + price_ticks_from_best: int + qty_fraction: float + ttl_ms: int + cancel_order_id: Optional[str] = None + reduce_only: bool = False + post_only: bool = False + metadata: Mapping[str, Any] = field(default_factory=dict) + + +@dataclass(frozen=True, slots=True) +class CounterpartyAction: + """Adversarially useful aggregate actions that alter book/fill outcomes.""" + role: AgentRole + kind: ActionKind + side: Optional[Side] + price_ticks_from_best: int + qty_fraction_of_top: float + toxicity: float = 0.0 + metadata: Mapping[str, Any] = field(default_factory=dict) + + +JointAction = Tuple[Any, ...] # (our_action, cp_action_1, cp_action_2, ...) + + +@dataclass(frozen=True, slots=True) +class PlannedPolicy: + """Output of the planner: distribution over actions + selected action.""" + actions: Tuple[FulfilmentAction, ...] + probabilities: Tuple[float, ...] + selected_action: FulfilmentAction + diagnostics: Mapping[str, Any] + + +@dataclass(frozen=True, slots=True) +class RiskDecision: + approved: bool + action: Optional[FulfilmentAction] + reason: str + adjusted: bool = False diff --git a/MALKHUT/malkhut/engine.py b/MALKHUT/malkhut/engine.py new file mode 100644 index 0000000..04f63fc --- /dev/null +++ b/MALKHUT/malkhut/engine.py @@ -0,0 +1,295 @@ +""" +FulfilmentEngine — hot-path orchestrator. + +Input: latest canonical MarketWorldState. +Output: exchange order action or no-op. + +Cadence: every 100 ms, or on book/fill/position/intent/kill-switch update. + +NEVER runs CMA-ES. Only loads frozen PolicySnapshot. + +ASEx integration: + All mutable state mutations go through ASExGuardedState + ASExWorker. + Validate-before-mutate semantics guarantee no races, no corrupted state. + One worker thread per state object. No locks. +""" +from __future__ import annotations + +import hashlib +import time +from typing import Callable, Optional + +from malkhut.state import ( + FulfilmentPolicyParams, + MarketWorldState, + HOT_PATH_BUDGET_MS, +) +from malkhut.actions import PlannedPolicy, RiskDecision +from malkhut.planner.sm_mcts import DecoupledUCBPlanner +from malkhut.risk.gate import RiskGate +from malkhut.venue.bingx.adapter import BingXVenueAdapter +from malkhut.counterparties import CounterpartyPolicy, default_counterparty_ecology +from malkhut.cwm import MinimalCryptoLOBCWM, CodeWorldModel +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.registry import PolicyRegistry, PolicyStage + + +class FulfilmentEngine: + """ + Hot-path orchestration with ASEx validate-before-mutate semantics. + + All mutable state mutations go through ASEx workers: + - FulfilmentWorker: book/account/intent/policy state + - RiskWorker: risk gate decisions + kill switch + + ASEx guarantees: + - One writer thread per state (no races) + - Validate before apply (no corrupted state) + - No locks, no GC pressure + + Policy loading: + - params_provider: called each tick to get current policy + - registry: loads ACTIVE policy from CH + - HOT_RELOAD_POLICY: hot-swap via control plane + """ + + def __init__( + self, + params_provider: Optional[Callable[[], FulfilmentPolicyParams]] = None, + counterparties: Optional[tuple[CounterpartyPolicy, ...]] = None, + venue: Optional[BingXVenueAdapter] = None, + zinc: Optional[MalkhutZincPlane] = None, + control_plane: Optional[MalkhutControlPlane] = None, + store: Optional[MalkhutCHStore] = None, + registry: Optional[PolicyRegistry] = None, + ) -> None: + self.cwm: CodeWorldModel = MinimalCryptoLOBCWM() + self.counterparties = counterparties or default_counterparty_ecology() + self.planner = DecoupledUCBPlanner( + cwm=self.cwm, + counterparties=self.counterparties, + ) + self.venue = venue or BingXVenueAdapter() + self.zinc = zinc + self.control_plane = control_plane + self.store = store + self._active = True + + # Policy management + self._registry = registry or PolicyRegistry(store=store) + self._params_provider = params_provider + self._current_params: Optional[FulfilmentPolicyParams] = None + + # ASEx workers — serialised state mutations + self._fulfilment_worker = FulfilmentWorker() + self._risk_worker = RiskWorker() + + # Try to load active policy from registry + active = self._registry.load_active() + if active: + self._current_params = active + + @property + def params_provider(self) -> Callable[[], FulfilmentPolicyParams]: + """Get current policy. Falls back to registry → provider → default.""" + def _provide() -> FulfilmentPolicyParams: + # 1. Check if we have a cached params + if self._current_params is not None: + return self._current_params + # 2. Try registry + active = self._registry.load_active() + if active: + self._current_params = active + return active + # 3. Try provider + if self._params_provider: + return self._params_provider() + # 4. Default + from malkhut.state import FulfilmentPolicyParams + return FulfilmentPolicyParams( + version="default", 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, + ) + return _provide + + @params_provider.setter + def params_provider(self, value: Callable[[], FulfilmentPolicyParams]) -> None: + self._params_provider = value + + def hot_reload_policy(self, params: FulfilmentPolicyParams) -> None: + """Hot-reload a new policy (e.g. from control plane or registry).""" + self._current_params = params + if hasattr(self, '_fulfilment_worker'): + self._fulfilment_worker.reload_policy(params) + + def on_state(self, state: MarketWorldState) -> None: + """Process one state update through the full ASEx-guarded pipeline.""" + # 1. Check control plane for commands + self._process_control_plane() + + if not self._active: + return + + # 2. Load current policy parameters + params = self.params_provider() + + # 3. Plan + t0 = time.perf_counter_ns() + planned = self.planner.plan( + root_state=state, + params=params, + budget_ms=HOT_PATH_BUDGET_MS // 2, + ) + plan_ns = time.perf_counter_ns() - t0 + + # 4. Risk gate via ASEx + risk_check = RiskCheck(state=state, planned=planned, params=params) + risk_future = self._risk_worker._worker.mutate(risk_check) + decision = risk_future.result(timeout=5.0) + + # 5. Publish to Zinc shared memory + if self.zinc: + self._publish_to_zinc(state, planned, decision, plan_ns) + + # 6. Log to ClickHouse + if self.store: + self._persist_decision(state, planned, decision, plan_ns, params) + + # 7. Execute + self.venue.execute(state, decision) + + def _process_control_plane(self) -> None: + """Read and process commands from the CONTROL_PLANE region.""" + if not self.control_plane: + return + cmd = self.control_plane.read_command(timeout_ms=5) + if cmd is None: + return + + ts = time.time_ns() + + if cmd.command == ControlCommand.STOP.value: + self._active = False + if self.control_plane: + self.control_plane.publish_ack(cmd.command, ts, "acknowledged", "engine_stopped") + elif cmd.command == ControlCommand.START.value: + self._active = True + if self.control_plane: + self.control_plane.publish_ack(cmd.command, ts, "acknowledged", "engine_started") + elif cmd.command == ControlCommand.EMERGENCY_STOP.value: + self._active = False + if self.control_plane: + self.control_plane.publish_ack(cmd.command, ts, "acknowledged", "emergency_stop") + elif cmd.command == ControlCommand.STATUS_REQUEST.value: + status = "active" if self._active else "inactive" + policy_ver = self._current_params.version if self._current_params else "none" + if self.control_plane: + self.control_plane.publish_ack(cmd.command, ts, "status", + f"{status};policy={policy_ver}") + elif cmd.command == ControlCommand.HOT_RELOAD_POLICY.value: + # Load policy from registry by version (from params.policy_version) + version = cmd.params.get("policy_version", "") + if version: + record = self._registry.get_record(version) + if record and record.stage == PolicyStage.ACTIVE: + self.hot_reload_policy(record.params) + if self.control_plane: + self.control_plane.publish_ack(cmd.command, ts, "ok", + f"reloaded_{version}") + else: + if self.control_plane: + self.control_plane.publish_ack(cmd.command, ts, "error", + f"policy_{version}_not_active") + else: + # Reload from registry (latest ACTIVE) + active = self._registry.load_active() + if active: + self.hot_reload_policy(active) + if self.control_plane: + self.control_plane.publish_ack(cmd.command, ts, "ok", + f"reloaded_{active.version}") + else: + if self.control_plane: + self.control_plane.publish_ack(cmd.command, ts, "error", + "no_active_policy") + + def _publish_to_zinc( + self, state: MarketWorldState, planned: PlannedPolicy, + decision: RiskDecision, plan_ns: int, + ) -> None: + """Publish fulfilment output to Zinc shared memory.""" + self.zinc.publish_fulfilment({ + "ts_ns": state.ts_ns, + "symbol": state.venue.symbol, + "selected_action": str(decision.action.kind.value) if decision.action else "NONE", + "approved": decision.approved, + "risk_reason": decision.reason, + "plan_latency_ns": plan_ns, + "policy_version": "live", + "root_entropy": planned.diagnostics.get("entropy", 0.0), + "sims": planned.diagnostics.get("sims", 0), + }) + + self.zinc.publish_risk({ + "ts_ns": state.ts_ns, + "symbol": state.venue.symbol, + "approved": decision.approved, + "reason": decision.reason, + }) + + def _persist_decision( + self, state: MarketWorldState, planned: PlannedPolicy, + decision: RiskDecision, plan_ns: int, params: FulfilmentPolicyParams, + ) -> None: + """Log decision to ClickHouse.""" + state_hash = hashlib.sha256( + f"{state.ts_ns}:{state.venue.symbol}".encode() + ).hexdigest()[:16] + + self.store.store_fulfilment_decision( + ts_ns=state.ts_ns, + exchange=state.venue.exchange, + symbol=state.venue.symbol, + intent_id=state.intent.intent_id if state.intent else "", + state_hash=state_hash, + selected_action=str(decision.action.kind.value) if decision.action else "NONE", + root_distribution=str(planned.probabilities), + risk_decision=f"{decision.approved}:{decision.reason}", + policy_version=params.version, + latency_ms=plan_ns / 1_000_000.0, + ) + + def close(self) -> None: + """Shut down ASEx workers and clean up resources.""" + self._active = False + self._fulfilment_worker.close() + self._risk_worker.close() + if self.zinc: + self.zinc.close_all() + if self.control_plane: + self.control_plane.close() + + @property + def fulfilment_worker(self) -> FulfilmentWorker: + return self._fulfilment_worker + + @property + def risk_worker(self) -> RiskWorker: + return self._risk_worker diff --git a/MALKHUT/malkhut/features.py b/MALKHUT/malkhut/features.py new file mode 100644 index 0000000..60694db --- /dev/null +++ b/MALKHUT/malkhut/features.py @@ -0,0 +1,54 @@ +""" +Feature extraction for planner and reward functions. + +Rule: every human-obvious feature is allowed, but the CMA-ES optimiser +must be allowed to discover non-obvious interactions (queue churn, time +since MFE, recovery velocity, cross-venue lead, etc.). +""" +from __future__ import annotations + +import math +from dataclasses import dataclass +from typing import Mapping, Protocol + +from malkhut.state import MarketWorldState + + +@dataclass(frozen=True, slots=True) +class FeatureVector: + values: Mapping[str, float] + + +class FeatureExtractor(Protocol): + def extract(self, state: MarketWorldState) -> FeatureVector: ... + + +class DefaultFeatureExtractor: + def extract(self, state: MarketWorldState) -> FeatureVector: + b = state.book + bid_qty = sum(x.qty for x in b.bids[:5]) + ask_qty = sum(x.qty for x in b.asks[:5]) + imbalance = (bid_qty - ask_qty) / max(bid_qty + ask_qty, 1e-12) + + path = state.trade_path + + values = { + "mid": b.mid if b.bids and b.asks else 0.0, + "spread_bps": b.spread_bps if b.bids and b.asks else 0.0, + "top5_imbalance": imbalance, + "funding_bps": state.funding_bps or 0.0, + "volatility_state": state.volatility_state or 0.0, + "pnl_bps": path.pnl_bps if path else 0.0, + "mae_bps": path.mae_bps if path else 0.0, + "mfe_bps": path.mfe_bps if path else 0.0, + "distance_from_mfe_bps": path.distance_from_mfe_bps if path else 0.0, + "seconds_held": path.seconds_held if path else 0.0, + "time_in_loss_s": path.time_in_loss_s if path else 0.0, + "time_since_deep_mae_s": path.time_since_deep_mae_s if path else 0.0, + "recovery_velocity_bps_per_s": path.recovery_velocity_bps_per_s if path else 0.0, + "adverse_velocity_bps_per_s": path.adverse_velocity_bps_per_s if path else 0.0, + "orderflow_toxicity": path.orderflow_toxicity if path else 0.0, + "queue_churn_score": path.queue_churn_score if path else 0.0, + "cross_venue_lead_score": path.cross_venue_lead_score if path else 0.0, + } + return FeatureVector(values=values) diff --git a/MALKHUT/malkhut/state.py b/MALKHUT/malkhut/state.py new file mode 100644 index 0000000..bc8d992 --- /dev/null +++ b/MALKHUT/malkhut/state.py @@ -0,0 +1,333 @@ +""" +MALKHUT canonical data model. + +All state objects are frozen+slots for: + - deterministic tree search (immutable snapshots) + - GraalVM compatibility (no mutable default hell) + - lock-free shared memory (readers never see partial writes) +""" +from __future__ import annotations + +from dataclasses import dataclass, field +from enum import Enum +from typing import Any, Dict, Mapping, Optional, Sequence, Tuple +import math + + +# ============================================================================== +# Enums +# ============================================================================== + +class Side(str, Enum): + BUY = "BUY" + SELL = "SELL" + + +class OrderType(str, Enum): + LIMIT = "LIMIT" + MARKET = "MARKET" + POST_ONLY = "POST_ONLY" + IOC = "IOC" + FOK = "FOK" + REDUCE_ONLY_LIMIT = "REDUCE_ONLY_LIMIT" + REDUCE_ONLY_MARKET = "REDUCE_ONLY_MARKET" + + +class ActionKind(str, Enum): + NOOP = "NOOP" + PLACE = "PLACE" + CANCEL = "CANCEL" + CANCEL_REPLACE = "CANCEL_REPLACE" + CROSS_SPREAD = "CROSS_SPREAD" + REDUCE = "REDUCE" + FULL_EXIT = "FULL_EXIT" + MOVE_STOP = "MOVE_STOP" + MOVE_TAKE_PROFIT = "MOVE_TAKE_PROFIT" + THROTTLE = "THROTTLE" + + +class IntentKind(str, Enum): + ENTER_LONG = "ENTER_LONG" + ENTER_SHORT = "ENTER_SHORT" + ADD_LONG = "ADD_LONG" + ADD_SHORT = "ADD_SHORT" + REDUCE_LONG = "REDUCE_LONG" + REDUCE_SHORT = "REDUCE_SHORT" + EXIT_LONG = "EXIT_LONG" + EXIT_SHORT = "EXIT_SHORT" + MAINTAIN = "MAINTAIN" + + +class AgentRole(str, Enum): + OUR_FULFILMENT = "OUR_FULFILMENT" + PASSIVE_MAKER = "PASSIVE_MAKER" + TOXIC_TAKER = "TOXIC_TAKER" + LATENCY_ARB = "LATENCY_ARB" + MOMENTUM_TAKER = "MOMENTUM_TAKER" + MEAN_REVERSION_TAKER = "MEAN_REVERSION_TAKER" + INVENTORY_MM = "INVENTORY_MM" + LIQUIDATION_FLOW = "LIQUIDATION_FLOW" + NOISE_TRADER = "NOISE_TRADER" + STALE_QUOTE_ATTACKER = "STALE_QUOTE_ATTACKER" + + +class Mode(str, Enum): + REPLAY_NO_IMPACT = "REPLAY_NO_IMPACT" + ENDOGENOUS_AGENT_SIM = "ENDOGENOUS_AGENT_SIM" + PAPER = "PAPER" + SHADOW_LIVE = "SHADOW_LIVE" + LIVE = "LIVE" + + +# ============================================================================== +# Core constants +# ============================================================================== + +HOT_PATH_BUDGET_MS: int = 100 +DEFAULT_PLANNER_BUDGET_MS: int = 25 +DEFAULT_TREE_DEPTH: int = 3 +DEFAULT_MAX_SIMS: int = 256 +DEFAULT_UCB_C: float = 1.41421356237 +DEFAULT_MIN_ROOT_POLICY_ENTROPY: float = 0.25 +DEFAULT_SELF_PLAY_POOL_MAX: int = 12 +DEFAULT_POLICY_PROMOTION_MIN_EDGE_BPS: float = 0.75 +DEFAULT_POLICY_PROMOTION_MIN_PVALUE: float = 0.05 + +MAX_ACCOUNT_LEVERAGE: float = 2.0 +MAX_EXCHANGE_LEVERAGE: float = 5.0 +MAX_SINGLE_ORDER_NOTIONAL_FRACTION: float = 0.05 +MAX_SYMBOL_NOTIONAL_FRACTION: float = 0.20 +MAX_CANCELS_PER_SYMBOL_PER_MINUTE: int = 90 +TAIL_QUANTILE: float = 0.05 + + +# ============================================================================== +# Frozen data model +# ============================================================================== + +@dataclass(frozen=True, slots=True) +class VenueRules: + exchange: str + symbol: str + tick_size: float + lot_size: float + min_qty: float + min_notional: float + maker_fee_bps: float + taker_fee_bps: float + post_only_supported: bool + reduce_only_supported: bool + max_orders_per_second: int + max_cancels_per_minute: int + + +@dataclass(frozen=True, slots=True) +class PriceLevel: + price: float + qty: float + + +@dataclass(frozen=True, slots=True) +class OrderBookState: + ts_ns: int + symbol: str + bids: Tuple[PriceLevel, ...] + asks: Tuple[PriceLevel, ...] + last_trade_price: Optional[float] = None + last_trade_qty: Optional[float] = None + last_trade_side: Optional[Side] = None + + @property + def best_bid(self) -> float: + return self.bids[0].price + + @property + def best_ask(self) -> float: + return self.asks[0].price + + @property + def mid(self) -> float: + return 0.5 * (self.best_bid + self.best_ask) + + @property + def spread(self) -> float: + return self.best_ask - self.best_bid + + @property + def spread_bps(self) -> float: + return 10_000.0 * self.spread / max(self.mid, 1e-12) + + +@dataclass(frozen=True, slots=True) +class PositionState: + symbol: str + qty: float + avg_entry: float + unrealized_pnl: float + realized_pnl: float + liquidation_price: Optional[float] + leverage: float + side: Optional[Side] + + +@dataclass(frozen=True, slots=True) +class AccountState: + ts_ns: int + equity: float + wallet_balance: float + available_balance: float + margin_used: float + total_notional: float + positions: Mapping[str, PositionState] = field(default_factory=dict) + + +@dataclass(frozen=True, slots=True) +class OpenOrderState: + client_order_id: str + venue_order_id: Optional[str] + symbol: str + side: Side + order_type: OrderType + price: Optional[float] + qty: float + remaining_qty: float + queue_ahead_estimate: Optional[float] + created_ts_ns: int + last_update_ts_ns: int + reduce_only: bool = False + post_only: bool = False + + +@dataclass(frozen=True, slots=True) +class TradePathState: + """In-trade path encoding for path-aware SL/TP.""" + symbol: str + side: Side + entry_ts_ns: int + now_ts_ns: int + bars_held: int + seconds_held: float + + pnl_bps: float + mae_bps: float + mfe_bps: float + distance_from_mfe_bps: float + distance_from_entry_bps: float + + time_to_mfe_s: float + time_in_loss_s: float + time_in_profit_s: float + time_since_last_profit_s: float + time_since_deep_mae_s: float + + loss_to_profit_transitions: int + deep_loss_recoveries: int + failed_recovery_count: int + recovery_velocity_bps_per_s: float + adverse_velocity_bps_per_s: float + + dolphin_regime_score: float + jericho_signal_strength: float + volatility_bps: float + orderflow_toxicity: float + queue_churn_score: float + book_imbalance: float + cross_venue_lead_score: float + + +@dataclass(frozen=True, slots=True) +class ExecutionIntent: + intent_id: str + ts_ns: int + symbol: str + kind: IntentKind + target_qty: float + max_notional: float + urgency: float + alpha_horizon_s: float + alpha_bps: float + max_slippage_bps: float + prefer_maker: bool + reduce_only: bool + ttl_s: float + reason: str + + +@dataclass(frozen=True, slots=True) +class MarketWorldState: + """Complete CWM root state. Immutable for safe tree search.""" + ts_ns: int + mode: Mode + venue: VenueRules + book: OrderBookState + account: AccountState + open_orders: Tuple[OpenOrderState, ...] = () + trade_path: Optional[TradePathState] = None + intent: Optional[ExecutionIntent] = None + + funding_bps: Optional[float] = None + volatility_state: Optional[float] = None + market_regime: Optional[str] = None + + feed_latency_ms: float = 0.0 + order_latency_ms: float = 0.0 + rng_seed: int = 0 + + +@dataclass(frozen=True, slots=True) +class FulfilmentPolicyParams: + """ + Frozen parameter set loaded by the live planner. + CMA-ES tunes this object offline. + """ + version: str + + # Planner + ucb_c: float + max_sims: int + max_depth: int + rollout_depth: int + root_temperature: float + min_root_entropy: float + + # Quote menu + quote_offsets_ticks: Tuple[int, ...] + quote_size_fractions: Tuple[float, ...] + passive_ttl_ms: int + aggressive_ttl_ms: int + + # Maker/taker thresholds + maker_edge_min_bps: float + cross_spread_edge_min_bps: float + adverse_toxicity_cancel_threshold: float + queue_churn_cancel_threshold: float + + # SL/TP/path risk + mae_tail_cut_bps: float + mfe_giveback_cut_fraction: float + max_time_in_loss_s: float + failed_recovery_cut_count: int + recovery_velocity_min_bps_per_s: float + + # Inventory/account + max_symbol_notional_fraction: float + max_single_order_notional_fraction: float + reduce_when_global_up_fraction: float + session_profit_lock_fraction: float + + # Reward weights + w_expected_pnl: float + w_fill_probability: float + w_adverse_selection: float + w_queue_priority: float + w_inventory_risk: float + w_tail_loss: float + w_fee_quality: float + w_time_decay: float + w_policy_entropy: float + + # Scenario robustness + robust_tail_weight: float + toxic_counterparty_weight: float + low_liquidity_weight: float + latency_stress_weight: float