296 lines
12 KiB
Python
296 lines
12 KiB
Python
|
|
"""
|
||
|
|
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
|