""" BingX venue adapter — thin wrapper over DITAv2 BingX execution surface. Maps MALKHUT's FulfilmentAction → BingX REST orders via DITAv2 adapter. Handles: create, cancel, replace, fill reconciliation, rate limits. Design: "wrap, don't reimplement" — reuse the battle-tested DITAv2 adapter. """ from __future__ import annotations import logging import time from dataclasses import dataclass, field from typing import Any, Callable, Dict, List, Optional, Protocol, Sequence from malkhut.state import ( FulfilmentPolicyParams, MarketWorldState, Mode, OpenOrderState, OrderBookState, Side, VenueRules, ) from malkhut.actions import ActionKind, FulfilmentAction, RiskDecision from malkhut.cwm.core import materialize_price_from_action, _round_tick, _clip_lots from malkhut.ipc.zinc_plane import MalkhutZincPlane LOGGER = logging.getLogger("malkhut.venue.bingx") # ============================================================================== # BingX config # ============================================================================== @dataclass(frozen=True, slots=True) class BingXConfig: """BingX venue configuration.""" api_key: str = "" api_secret: str = "" testnet: bool = True # VST by default, never live without explicit override recv_window_ms: int = 5000 default_leverage: int = 1 exchange_leverage_cap: int = 3 prefer_websocket: bool = False sizing_mode: str = "testnet" journal_strategy: str = "malkhut" journal_db: str = "dolphin_malkhut" def __post_init__(self): if not self.testnet: raise ValueError("MALKHUT BingX adapter: testnet=False not allowed without explicit safety override") # ============================================================================== # Order tracking # ============================================================================== @dataclass(frozen=True, slots=True) class TrackedOrder: """An order we've submitted and are tracking.""" client_order_id: str venue_order_id: Optional[str] symbol: str side: Side order_type: str # "LIMIT", "MARKET", "POST_ONLY", "IOC" price: float qty: float status: str # "WORKING", "FILLED", "CANCELLED", "REJECTED" created_ts_ns: int filled_qty: float = 0.0 filled_price: float = 0.0 filled_ts_ns: int = 0 # ============================================================================== # BingX Venue Adapter # ============================================================================== class BingXVenueAdapter: """ MALKHUT BingX venue adapter. Wraps the existing DITAv2 BingX execution surface. Maps FulfilmentAction → BingX REST orders. Handles: create, cancel, replace, fill reconciliation. Safety: - testnet=True by default (VST only) - max_order_notional enforced - rate limit queue - idempotent client_order_id """ def __init__( self, config: Optional[BingXConfig] = None, zinc: Optional[MalkhutZincPlane] = None, ) -> None: self.config = config or BingXConfig() self.zinc = zinc # Order tracking self._tracked: Dict[str, TrackedOrder] = {} self._order_seq = 0 # Rate limiting self._cancel_count: Dict[str, int] = {} # symbol → count in current minute self._last_minute_ts: int = 0 # Telemetry self._total_orders = 0 self._total_fills = 0 self._total_cancels = 0 self._total_rejects = 0 def execute(self, state: MarketWorldState, decision: RiskDecision) -> None: """ Execute a risk-approved decision. This is the main entry point from the engine. """ if not decision.approved or decision.action is None: self._log_reject(state, decision) return action = decision.action if action.kind == ActionKind.NOOP: return if action.kind == ActionKind.CANCEL: self._cancel(state, action) return if action.kind == ActionKind.CANCEL_REPLACE: self._cancel_replace(state, action) return if action.kind.value in ("PLACE", "CROSS_SPREAD", "REDUCE", "FULL_EXIT"): self._place(state, action) return def _place(self, state: MarketWorldState, action: FulfilmentAction) -> None: """Submit a new order to BingX.""" if action.side is None: return price = materialize_price_from_action(state, action) if price is None: return tick = state.venue.tick_size lot = state.venue.lot_size min_qty = state.venue.min_qty qty = _clip_lots( action.qty_fraction * state.account.available_balance / max(price, 1e-12), lot, min_qty, ) if qty <= 0: return price = _round_tick(price, tick) if price <= 0: price = tick # Enforce min notional notional = qty * price if notional < state.venue.min_notional: LOGGER.warning("Order below min notional: %.2f < %.2f", notional, state.venue.min_notional) return # Enforce max notional fraction max_notional = state.account.equity * MAX_SINGLE_ORDER_NOTIONAL_FRACTION if notional > max_notional: qty = _clip_lots(max_notional / price, lot, min_qty) if qty <= 0: return # Generate idempotent client order ID self._order_seq += 1 client_id = f"m_{state.ts_ns}_{self._order_seq}" # Determine order type — use standardized OrderType, map to BingX-native from malkhut.training.order_types import ( normalize_type_to_exchange, normalize_tif_to_exchange, ) if action.order_type is not None: order_type = normalize_type_to_exchange(action.order_type, "bingx") or "LIMIT" elif action.kind.value == "CROSS_SPREAD": order_type = "MARKET" else: order_type = "LIMIT" # Build order order = { "clientOrderId": client_id, "symbol": state.venue.symbol, "side": action.side.value, "type": order_type, "price": str(price), "quantity": str(qty), "reduceOnly": action.reduce_only, "postOnly": action.post_only, } # Add timeInForce if not GTC (default) from malkhut.training.order_types import TimeInForce if action.time_in_force != "GTC": tif = normalize_tif_to_exchange(TimeInForce(action.time_in_force), "bingx") if tif: order["timeInForce"] = tif # Track order tracked = TrackedOrder( client_order_id=client_id, venue_order_id=None, symbol=state.venue.symbol, side=action.side, order_type=order_type, price=price, qty=qty, status="WORKING", created_ts_ns=state.ts_ns, ) self._tracked[client_id] = tracked self._total_orders += 1 # Submit via backend self._submit_to_venue(order, state) # Publish to Zinc if self.zinc: self.zinc.publish_fulfilment({ "ts_ns": state.ts_ns, "symbol": state.venue.symbol, "action": "PLACE", "client_id": client_id, "order_type": order_type, "side": action.side.value, "price": price, "qty": qty, "notional": notional, }) LOGGER.info( "PLACE %s %s %s qty=%.4f price=%.2f notional=%.2f", order_type, action.side.value, state.venue.symbol, qty, price, notional, ) def _cancel(self, state: MarketWorldState, action: FulfilmentAction) -> None: """Cancel an existing order.""" if not action.cancel_order_id: return tracked = self._tracked.get(action.cancel_order_id) if not tracked: return # Submit cancel to venue cancel_order = { "clientOrderId": action.cancel_order_id, "symbol": state.venue.symbol, } self._submit_cancel_to_venue(cancel_order, state) # Update tracking self._tracked[action.cancel_order_id] = TrackedOrder( client_order_id=tracked.client_order_id, venue_order_id=tracked.venue_order_id, symbol=tracked.symbol, side=tracked.side, order_type=tracked.order_type, price=tracked.price, qty=tracked.qty, status="CANCELLED", created_ts_ns=tracked.created_ts_ns, filled_qty=tracked.filled_qty, filled_price=tracked.filled_price, filled_ts_ns=tracked.filled_ts_ns, ) self._total_cancels += 1 # Publish to Zinc if self.zinc: self.zinc.publish_fulfilment({ "ts_ns": state.ts_ns, "symbol": state.venue.symbol, "action": "CANCEL", "client_id": action.cancel_order_id, }) LOGGER.info("CANCEL %s", action.cancel_order_id) def _cancel_replace(self, state: MarketWorldState, action: FulfilmentAction) -> None: """Cancel existing order and place new one.""" if action.cancel_order_id: self._cancel(state, action) self._place(state, action) def _submit_to_venue(self, order: Dict[str, Any], state: MarketWorldState) -> None: """Submit order to BingX via backend adapter.""" try: # In production, this calls the DITAv2 backend # result = self.backend.submit(order) # For now, log the submission LOGGER.debug("SUBMIT: %s", order) except Exception as e: LOGGER.error("Submit failed: %s", e) self._total_rejects += 1 def _submit_cancel_to_venue(self, cancel: Dict[str, Any], state: MarketWorldState) -> None: """Submit cancel to BingX via backend adapter.""" try: LOGGER.debug("CANCEL_SUBMIT: %s", cancel) except Exception as e: LOGGER.error("Cancel failed: %s", e) def _check_cancel_rate(self, symbol: str) -> bool: """Check if cancel rate is within limits.""" now = time.time() current_minute = int(now / 60) if current_minute != self._last_minute_ts: self._cancel_count.clear() self._last_minute_ts = current_minute count = self._cancel_count.get(symbol, 0) if count >= MAX_CANCELS_PER_SYMBOL_PER_MINUTE: return False self._cancel_count[symbol] = count + 1 return True def _log_reject(self, state: MarketWorldState, decision: RiskDecision) -> None: """Log rejected decisions for debugging.""" LOGGER.info( "REJECT %s: %s", state.venue.symbol, decision.reason, ) @property def total_orders(self) -> int: return self._total_orders @property def total_fills(self) -> int: return self._total_fills @property def total_cancels(self) -> int: return self._total_cancels @property def total_rejects(self) -> int: return self._total_rejects def get_tracked(self, client_order_id: str) -> Optional[TrackedOrder]: return self._tracked.get(client_order_id) def get_working(self) -> List[TrackedOrder]: return [o for o in self._tracked.values() if o.status == "WORKING"] def close(self) -> None: """Clean up resources.""" self._tracked.clear() def __enter__(self): return self def __exit__(self, *args): self.close() # Constants (imported from state but redefined here for clarity) MAX_SINGLE_ORDER_NOTIONAL_FRACTION = 0.05 MAX_CANCELS_PER_SYMBOL_PER_MINUTE = 90