From 4aeadf1aae8036d8082f7a4a8b6fd97a975463cb Mon Sep 17 00:00:00 2001 From: Codex Date: Tue, 14 Jul 2026 18:37:41 +0200 Subject: [PATCH] =?UTF-8?q?docs:=20hftbacktest=20CWM=20integration=20desig?= =?UTF-8?q?n=20=E2=80=94=20drop-in=20LOB=20backend?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- MALKHUT/docs/HFTBACKTEST_CWM_INTEGRATION.md | 359 ++++++++++++++++++++ 1 file changed, 359 insertions(+) create mode 100644 MALKHUT/docs/HFTBACKTEST_CWM_INTEGRATION.md diff --git a/MALKHUT/docs/HFTBACKTEST_CWM_INTEGRATION.md b/MALKHUT/docs/HFTBACKTEST_CWM_INTEGRATION.md new file mode 100644 index 0000000..032c5d4 --- /dev/null +++ b/MALKHUT/docs/HFTBACKTEST_CWM_INTEGRATION.md @@ -0,0 +1,359 @@ +# 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 + +```python +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: + +```python +# 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: + +```python +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.