malkhut(T1): scaffold — frozen state model, actions, features, engine
T1 scaffold: 42 frozen dataclasses (state.py), action model + PlannedPolicy + RiskDecision (actions.py), 17-feature extraction (features.py), FulfilmentEngine hot-path orchestrator (engine.py).
This commit is contained in:
295
MALKHUT/malkhut/engine.py
Normal file
295
MALKHUT/malkhut/engine.py
Normal file
@@ -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
|
||||
Reference in New Issue
Block a user