""" 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