CWM (103): core mechanics, exhaustive edge cases, numba, exchange mechanics Replay (118): exhaustive verification, microstructure, trajectory Training (190): asset classification, phase0 extensive, pipeline, exhaustive DSL (102): v2 syntax, expanded, new features ASEx (33): validate-before-mutate, single-writer Planner (48): MCTS, alternatives, hooks Counterparties (19): 9 adversarial agent policies Clock (30): event-driven reactor BingX (28): venue adapter IPC (8): Zinc SHM Storage (9): ClickHouse Risk (4): hard invariants State (17): frozen dataclass invariants Integration: E2E, concurrency, sync/async seams, hypothesis, fuzz, adversarial
178 lines
6.5 KiB
Python
178 lines
6.5 KiB
Python
"""
|
|
Sync/async seam tests.
|
|
|
|
Tests the boundary between synchronous and asynchronous components:
|
|
- Zinc SHM (synchronous POSIX SHM) with async-like polling
|
|
- Control plane command processing in event loop
|
|
- Engine hot-path with timeout budget
|
|
"""
|
|
import threading
|
|
import time
|
|
import pytest
|
|
from malkhut.ipc.zinc_plane import MalkhutZincPlane, SharedRegionWriter, SharedRegionReader
|
|
from malkhut.ipc.control_plane import MalkhutControlPlane, ControlPlaneFrame
|
|
from malkhut.engine import FulfilmentEngine
|
|
from malkhut.state import (
|
|
AccountState, FulfilmentPolicyParams, MarketWorldState, Mode,
|
|
OrderBookState, PriceLevel, VenueRules,
|
|
)
|
|
|
|
|
|
def _params():
|
|
return FulfilmentPolicyParams(
|
|
version="seam", ucb_c=1.414, max_sims=32, max_depth=2,
|
|
rollout_depth=1, root_temperature=0.5, min_root_entropy=0.25,
|
|
quote_offsets_ticks=(0, 1), quote_size_fractions=(0.25,),
|
|
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,
|
|
)
|
|
|
|
|
|
def _state():
|
|
return MarketWorldState(
|
|
ts_ns=1_000_000_000, mode=Mode.PAPER,
|
|
venue=VenueRules(
|
|
exchange="bingx", symbol="BTCUSDT", tick_size=0.1, lot_size=0.001,
|
|
min_qty=0.001, min_notional=5.0, maker_fee_bps=-0.2, taker_fee_bps=0.5,
|
|
post_only_supported=True, reduce_only_supported=True,
|
|
max_orders_per_second=100, max_cancels_per_minute=120,
|
|
),
|
|
book=OrderBookState(
|
|
ts_ns=1_000_000_000, symbol="BTCUSDT",
|
|
bids=(PriceLevel(50000.0, 1.0),), asks=(PriceLevel(50001.0, 1.0),),
|
|
),
|
|
account=AccountState(
|
|
ts_ns=1_000_000_000, equity=10000.0, wallet_balance=10000.0,
|
|
available_balance=10000.0, margin_used=0.0, total_notional=0.0,
|
|
),
|
|
)
|
|
|
|
|
|
class TestZincSyncSeam:
|
|
def test_write_is_synchronous(self):
|
|
"""Write completes before returning."""
|
|
plane = MalkhutZincPlane(prefix="sync_write")
|
|
t0 = time.perf_counter_ns()
|
|
plane.publish_book({"ts": time.time_ns()})
|
|
elapsed_ms = (time.perf_counter_ns() - t0) / 1_000_000
|
|
assert elapsed_ms < 100 # should be very fast
|
|
plane.close_all()
|
|
|
|
def test_read_timeout_returns_within_budget(self):
|
|
"""Read with timeout returns within the timeout window."""
|
|
plane = MalkhutZincPlane(prefix="sync_read_timeout")
|
|
plane.publish_book({"test": True})
|
|
t0 = time.perf_counter_ns()
|
|
data, seq = plane.read_book(timeout_ms=50)
|
|
elapsed_ms = (time.perf_counter_ns() - t0) / 1_000_000
|
|
assert elapsed_ms < 200 # should complete within timeout
|
|
plane.close_all()
|
|
|
|
def test_write_read_roundtrip_latency(self):
|
|
"""Write-to-read latency is measurable and bounded."""
|
|
plane = MalkhutZincPlane(prefix="latency_test")
|
|
plane.publish_book({"ts": time.time_ns()})
|
|
t0 = time.perf_counter_ns()
|
|
data, seq = plane.read_book(timeout_ms=50)
|
|
latency_us = (time.perf_counter_ns() - t0) / 1_000
|
|
assert latency_us < 10_000 # < 10ms
|
|
plane.close_all()
|
|
|
|
|
|
class TestControlPlaneSyncSeam:
|
|
def test_command_processing_latency(self):
|
|
"""Control plane command write-to-read is fast."""
|
|
cp = MalkhutControlPlane()
|
|
frame = ControlPlaneFrame(
|
|
command="STATUS_REQUEST", ts_ns=time.time_ns(), source="test",
|
|
)
|
|
t0 = time.perf_counter_ns()
|
|
cp.publish_command(frame)
|
|
cmd = cp.read_command(timeout_ms=50)
|
|
elapsed_us = (time.perf_counter_ns() - t0) / 1_000
|
|
assert cmd is not None
|
|
assert elapsed_us < 10_000 # < 10ms
|
|
cp.close()
|
|
|
|
|
|
class TestEngineSeam:
|
|
def test_engine_hot_path_latency(self):
|
|
"""Engine on_state completes within budget."""
|
|
from malkhut.venue.bingx.adapter import BingXVenueAdapter
|
|
engine = FulfilmentEngine(
|
|
params_provider=_params,
|
|
venue=BingXVenueAdapter(),
|
|
)
|
|
t0 = time.perf_counter_ns()
|
|
engine.on_state(_state())
|
|
elapsed_ms = (time.perf_counter_ns() - t0) / 1_000_000
|
|
assert elapsed_ms < 200 # hot path budget
|
|
|
|
def test_engine_stop_start_via_control_plane(self):
|
|
"""Engine can be stopped and started via control plane."""
|
|
from malkhut.ipc.control_plane import MalkhutControlPlane, ControlPlaneFrame
|
|
cp = MalkhutControlPlane()
|
|
engine = FulfilmentEngine(
|
|
params_provider=_params,
|
|
control_plane=cp,
|
|
)
|
|
|
|
# Initially active
|
|
assert engine._active
|
|
|
|
# Stop via control plane
|
|
cp.publish_command(ControlPlaneFrame(
|
|
command="STOP", ts_ns=time.time_ns(), source="test",
|
|
))
|
|
engine.on_state(_state())
|
|
assert not engine._active
|
|
|
|
# Start via control plane
|
|
cp.publish_command(ControlPlaneFrame(
|
|
command="START", ts_ns=time.time_ns(), source="test",
|
|
))
|
|
engine.on_state(_state())
|
|
assert engine._active
|
|
|
|
cp.close()
|
|
|
|
|
|
class TestCrossThreadSync:
|
|
def test_writer_in_thread_reader_in_main(self):
|
|
"""Writer in background thread, reader in main thread."""
|
|
plane = MalkhutZincPlane(prefix="cross_thread")
|
|
written = threading.Event()
|
|
|
|
def bg_writer():
|
|
for i in range(5):
|
|
plane.publish_book({"thread_seq": i})
|
|
time.sleep(0.01)
|
|
written.set()
|
|
|
|
t = threading.Thread(target=bg_writer)
|
|
t.start()
|
|
|
|
results = []
|
|
for _ in range(10):
|
|
try:
|
|
data, seq = plane.read_book(timeout_ms=20)
|
|
results.append(data)
|
|
except Exception:
|
|
pass
|
|
time.sleep(0.005)
|
|
|
|
t.join(timeout=5)
|
|
plane.close_all()
|
|
assert len(results) > 0
|