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
This commit is contained in:
359
MALKHUT/docs/HFTBACKTEST_CWM_INTEGRATION.md
Normal file
359
MALKHUT/docs/HFTBACKTEST_CWM_INTEGRATION.md
Normal file
@@ -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.
|
||||
Reference in New Issue
Block a user