Files
sentiment-engine/MALKHUT/malkhut/engine.py

296 lines
12 KiB
Python
Raw Normal View History

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