Files
sentiment-engine/MALKHUT/docs/HFTBACKTEST_CWM_INTEGRATION.md
Codex 4aeadf1aae docs: hftbacktest CWM integration design — drop-in LOB backend
Architecture:
  hftbacktest replaces _fill_from_levels() + manual book updates
  Everything above transition() stays the same

What changes:
  - HftBacktestCWM: new class implementing CodeWorldModel protocol
  - transition(): submit/cancel via hftbacktest, convert state back
  - fill model: ProbQueueModel (queue-position-aware)
  - latency: interpolated from historical data

What does NOT change:
  - Planner (SM-MCTS, EXP3, Thompson, etc.)
  - Counterparty ecology (ToxicTaker, PassiveMaker, etc.)
  - Risk gate (all 6 checks)
  - CMA-ES trainer
  - PerformanceMatrix
  - Reward function (same PnL + adverse selection + risk)
  - Action menu, FulfilmentAction, ScenarioFactory
  - ALL existing tests

Integration: 3 steps, zero core changes:
  1. Create malkhut/cwm/hft_cwm.py
  2. Swap cwm_factory in PolicyEvaluator
  3. Done
2026-07-14 18:37:41 +02:00

16 KiB
Raw Blame History

hftbacktest CWM Integration Design

Goal: Use hftbacktest as the simulated exchange/OB engine underneath MALKHUT's CWM, while keeping the entire game-theoretic layer (planner, counterparty ecology, CMA-ES, risk gate, PerformanceMatrix) unchanged.

Principle: hftbacktest replaces _fill_from_levels() + manual book updates. Everything above transition() stays the same.


Architecture: What Changes, What Doesn't

                    MALKHUT (unchanged)
┌──────────────────────────────────────────────────────────────┐
│  Planner (DecoupledUCBPlanner / EXP3 / Thompson / ...)       │
│  CounterpartyEcology (ToxicTaker / PassiveMaker / ...)       │
│  RiskGate (kill_switch / self_trade / leverage / ...)        │
│  CMA-ES Trainer + PerformanceMatrix + StrategySelector       │
│  FulfilmentAction (order_type × time_in_force × post_only)   │
└──────────────┬───────────────────────────────────────────────┘
               │ calls transition(state, joint_action)
               ▼
┌──────────────────────────────────────────────────────────────┐
│  CWM Protocol: transition() / reward() / terminal()          │
│  ┌────────────────────────────────────────────────────────┐  │
│  │  HftBacktestCWM  (NEW — replaces MinimalCryptoLOBCWM)  │  │
│  │                                                        │  │
│  │  transition() → hftbacktest submit/cancel + elapse     │  │
│  │  reward()     → MALKHUT reward function (unchanged)    │  │
│  │  terminal()   → unchanged                              │  │
│  └────────────────────────────────────────────────────────┘  │
└──────────────┬───────────────────────────────────────────────┘
               │ internally calls
               ▼
┌──────────────────────────────────────────────────────────────┐
│  hftbacktest HashMapMarketDepthBacktest                     │
│  (Rust-backed, event-driven LOB simulation)                 │
│                                                              │
│  .submit_buy_order()   ← our PLACE/CROSS_SPREAD             │
│  .submit_sell_order()  ← our PLACE/CROSS_SPREAD             │
│  .cancel()             ← our CANCEL                         │
│  .elapse(nanoseconds)  ← time progression                   │
│  .depth()              ← current book snapshot              │
│  .position()           ← our current position               │
│                                                              │
│  Features:                                                   │
│  - ProbQueueModel: probabilistic fill based on queue pos    │
│  - Interpolated latency: exchange + local event ordering    │
│  - Partial fills: order fills across multiple levels        │
│  - Fee models: flat_per_trade or trading_value              │
│  - Tick/lot: enforced by the engine                         │
└──────────────────────────────────────────────────────────────┘

The Bridge: HftBacktestCWM

class HftBacktestCWM:
    """CWM backed by hftbacktest's event-driven LOB engine.
    
    Implements the same CodeWorldModel protocol as MinimalCryptoLOBCWM.
    Drop-in replacement: same transition() / reward() / terminal() API.
    """

    def __init__(
        self,
        symbol: str = "BTCUSDT",
        tick_size: float = 0.1,
        lot_size: float = 0.001,
        maker_fee_bps: float = 2.0,
        taker_fee_bps: float = 5.0,
        latency_ns: int = 100_000_000,  # 100ms order latency
        data: Optional[np.ndarray] = None,  # pre-loaded L2 event data
    ):
        import hftbacktest as hbt
        
        asset = (hbt.BacktestAsset()
            .linear_asset(1.0)  # linear (not inverse) perp
            .tick_size(tick_size)
            .lot_size(lot_size)
            .flat_per_trade_fee_model(maker_fee_bps / 10_000, 
                                       taker_fee_bps / 10_000)
            .constant_order_latency(latency_ns, latency_ns)
            .power_prob_queue_model(3)  # queue position model
            .partial_fill_exchange()
        )
        
        if data is not None:
            asset.add_data(data)
        
        self.hbt = hbt.build_hashmap_backtest([asset])
        self._symbol = symbol
        self._tick_size = tick_size
        self._lot_size = lot_size
        self._order_id_seq = 0
        self._pending_fills = []  # filled orders awaiting retrieval

    def transition(
        self,
        state: MarketWorldState,
        joint_action: JointAction,
    ) -> MarketWorldState:
        our_action = joint_action[0]
        counterparty_actions = joint_action[1:]
        
        # 1. Process our action through hftbacktest
        if isinstance(our_action, FulfilmentAction):
            self._process_our_action(our_action, state)
        
        # 2. Process counterparty actions through hftbacktest
        for cp in counterparty_actions:
            if isinstance(cp, CounterpartyAction):
                self._process_counterparty(cp, state)
        
        # 3. Elapse time (advance the engine by one tick)
        self.hbt.elapse(1_000_000)  # 1ms
        
        # 4. Wait for order responses
        self.hbt.wait_next_feed()
        self.hbt.wait_order_response()
        
        # 5. Convert hftbacktest state → MALKHUT MarketWorldState
        return self._build_next_state(state, our_action)
    
    def _process_our_action(self, action: FulfilmentAction, state: MarketWorldState):
        """Convert MALKHUT FulfilmentAction → hftbacktest order submission."""
        import hftbacktest as hbt
        
        if action.kind.value in ("PLACE", "CANCEL_REPLACE"):
            price = materialize_price_from_action(state, action)
            if price is None:
                return
            qty = action.qty_fraction * state.account.available_balance / max(price, 1e-12)
            qty = _round_lot(qty, self._lot_size)
            if qty <= 0:
                return
            
            self._order_id_seq += 1
            if action.side == Side.BUY:
                self.hbt.submit_buy_order(
                    self._order_id_seq, qty, price,
                    hbt.Trigger.GTC,
                )
            else:
                self.hbt.submit_sell_order(
                    self._order_id_seq, qty, price,
                    hbt.Trigger.GTC,
                )
        
        elif action.kind.value == "CROSS_SPREAD":
            # Aggressive fill: submit at best available
            price = materialize_price_from_action(state, action)
            if price is None:
                return
            qty = action.qty_fraction * state.account.available_balance / max(price, 1e-12)
            qty = _round_lot(qty, self._lot_size)
            if qty <= 0:
                return
            
            self._order_id_seq += 1
            # Submit IOC-like (aggressive limit at market price)
            if action.side == Side.BUY:
                self.hbt.submit_buy_order(
                    self._order_id_seq, qty, state.book.best_ask,
                    hbt.Trigger.IOC,
                )
            else:
                self.hbt.submit_sell_order(
                    self._order_id_seq, qty, state.book.best_bid,
                    hbt.Trigger.IOC,
                )
        
        elif action.kind.value == "CANCEL":
            if action.cancel_order_id:
                oid = self._parse_order_id(action.cancel_order_id)
                self.hbt.cancel(oid)
    
    def _process_counterparty(self, cp: CounterpartyAction, state: MarketWorldState):
        """Counterparty actions hit the hftbacktest book as external events."""
        if cp.kind.value == "CROSS_SPREAD" and cp.side:
            # Counterparty crosses spread → inject as external trade
            qty = cp.qty_fraction_of_top * state.account.available_balance / max(
                state.book.mid if state.book.bids and state.book.asks else 1.0, 1e-12)
            price = state.book.best_ask if cp.side == Side.BUY else state.book.best_bid
            # hftbacktest handles this via feed events (external trades)
            # For simplicity, we submit as IOC from "other" side
            self._order_id_seq += 1
            if cp.side == Side.BUY:
                self.hbt.submit_sell_order(
                    self._order_id_seq, qty, price, hbt.Trigger.IOC,
                )
            else:
                self.hbt.submit_buy_order(
                    self._order_id_seq, qty, price, hbt.Trigger.IOC,
                )
    
    def _build_next_state(
        self,
        prev_state: MarketWorldState,
        action: FulfilmentAction,
    ) -> MarketWorldState:
        """Convert hftbacktest engine state → MALKHUT MarketWorldState."""
        
        # Get current position from hftbacktest
        hbt_pos = self.hbt.position(0)  # asset index 0
        
        # Get current book depth
        bid_depth = self.hbt.depth(0, is_ask=False)  # bid levels
        ask_depth = self.hbt.depth(0, is_ask=True)   # ask levels
        
        # Convert to MALKHUT OrderBookState
        bids = tuple(
            PriceLevel(float(level.px), float(level.qty))
            for level in bid_depth[:20]  # top 20 levels
            if level.qty > 0
        )
        asks = tuple(
            PriceLevel(float(level.px), float(level.qty))
            for level in ask_depth[:20]
            if level.qty > 0
        )
        
        book = OrderBookState(
            ts_ns=prev_state.ts_ns + 1_000_000,
            symbol=self._symbol,
            bids=bids or (PriceLevel(0.0, 0.0),),
            asks=asks or (PriceLevel(0.0, 0.0),),
        )
        
        # Convert position
        pos_qty = float(hbt_pos.qty)
        pos_avg = float(hbt_pos.avg_entry_price) if pos_qty != 0 else 0.0
        
        # ... (equity, available_balance, path_state calculation same as current CWM)
        
        return MarketWorldState(
            ts_ns=prev_state.ts_ns + 1_000_000,
            mode=prev_state.mode,
            venue=prev_state.venue,
            book=book,
            account=new_account,
            open_orders=(),  # hftbacktest tracks internally
            trade_path=new_trade_path,
            intent=prev_state.intent,
        )
    
    def reward(self, prev_state, action, next_state, params):
        """Same reward function as current CWM — unchanged."""
        return compute_reward_vectorized(...)
    
    def terminal(self, state, depth):
        """Same terminal check — unchanged."""
        return depth <= 0

Data Flow: How Actions Become Fills

Step 1: Planner calls plan(state, params) → PlannedPolicy
        selected_action = FulfilmentAction(PLACE, BUY, LIMIT, offset=5, tif=IOC)

Step 2: CMA-ES calls transition(state, (our_action, cp1, cp2, cp3))

Step 3: HftBacktestCWM.transition():
        a. submit_buy_order(id=42, qty=0.01, price=63999.5, IOC)
        b. Counterparty ToxicTaker: submit_sell_order(id=43, qty=0.005, IOC)
        c. hbt.elapse(1ms) → engine processes events
        d. hbt.wait_order_response() → fills collected
        e. _build_next_state() → MarketWorldState with updated book/position

Step 4: CMA-ES calls reward(prev, action, next, params)
        → Same reward function (unchanged)

Step 5: Repeat for next step

What We Get vs Current CWM

Feature Current CWM hftbacktest CWM
Fill model Deterministic level consumption Probabilistic queue position (PowerProbQueue)
Queue position Estimated (qty * 0.5) Modeled from order arrival/cancel dynamics
Latency Instant fill Interpolated from historical (100ms exchange latency)
Partial fills Yes (level-by-level) Yes (queue-aware)
Market impact Simple 0.5 * fraction Implicit in book consumption + refill
Fee model Manual calculation Built-in (flat_per_trade)
Counterparty fills External trade injection Same (IOC orders from other side)
Reward function MALKHUT custom UNCHANGED — same PnL + adverse selection + risk
Path state MALKHUT MAE/MFE UNCHANGED
Risk gate MALKHUT RiskGate UNCHANGED
Planner MALKHUT SM-MCTS UNCHANGED

Data Requirement

hftbacktest needs L2 depth data in its event array format:

# Event array dtype:
# (ev, exch_ts, local_ts, px, qty, order_id, ival, fval)
# ev: event type (1=depth, 2=trade, etc.)
# exch_ts: exchange timestamp (nanoseconds)
# local_ts: local receive timestamp (nanoseconds)
# px: price (float64)
# qty: quantity (float64)

data = hbt.Recorder.data("BTCUSDT", "2026-07-01")

Sources:

  • Tardis.dev (tardis.dev) — historical L2 data for Binance, Bybit, etc.
  • Binance data portal — free daily L2 snapshots
  • Live recording — hftbacktest has LiveInstrument for real-time capture

For our current use case (behavior-driven simulation), we can also SYNTHESIZE L2 data from our AssetBehavior profiles:

def synthesize_l2_data(behavior: AssetBehavior, duration_ns: int) -> np.ndarray:
    """Generate synthetic L2 events matching the asset's behavior profile."""
    events = []
    mid = behavior.reference_price
    for t in range(0, duration_ns, 1_000_000):  # 1ms steps
        # Generate depth events from power-law profile
        for d_bps in range(1, 100):
            depth_usd = behavior.depth_at_bps(d_bps)
            price = mid * (1 + d_bps / 10_000)
            events.append(make_depth_event(t, price, depth_usd / mid))
        # Generate trade events from flow profile
        n_trades = int(behavior.flow.orders_per_sec_normal / 1000)
        for _ in range(n_trades):
            trade_price = mid * (1 + random.gauss(0, behavior.vol.annualized_normal / 100))
            events.append(make_trade_event(t, trade_price, behavior.flow.avg_trade_usd / trade_price))
    return np.array(events, dtype=EVENT_ARRAY)

Integration Steps (no code changes to MALKHUT core)

  1. Create malkhut/cwm/hft_cwm.py — HftBacktestCWM class implementing CodeWorldModel protocol (transition/reward/terminal).

  2. Wire create_planner() to accept CWM class — already supports this: create_planner("sm_mcts", cwm=HftBacktestCWM(...), ...)

  3. Update PolicyEvaluator.cwm_factory — swap MinimalCryptoLOBCWM() with HftBacktestCWM(symbol=..., data=...).

  4. No changes to: planner, counterparty ecology, CMA-ES, risk gate, PerformanceMatrix, ScenarioFactory, action menu, or any test.

Why This Is Safe

The CWM is a leaf dependency — nothing depends ON it except the evaluator and the planner, both of which use it through the CodeWorldModel protocol. Swapping the implementation behind that protocol is a textbook Strategy pattern. The planner doesn't know or care whether the book is synthesized or hftbacktest. The reward function is pure math on (prev_state, action, next_state) — identical regardless of how next_state was computed.