use crate::compat::sleep;
use crate::performance_metrics::{
BacktestClosedTrade, BacktestEquityPoint, PerformanceBenchmark, PerformanceMetrics,
PerformanceMetricsSnapshot, PerformanceTracker,
};
use crate::poll_guard::{PollGuard, PollOutcome};
use bot_core::{
CancelAll, CancelOrder, ClientOrderId, Command, Event, Exchange, ExchangeError, ExchangeHealth,
ExchangeInstance, Fill, InstrumentId, OrderAcceptedEvent, OrderCanceledEvent,
OrderCompletedEvent, OrderFilledEvent, OrderInput, OrderRejectedEvent, PlaceOrder,
PlaceOrderResult, Position, Qty, Quote, QuoteEvent,
};
use futures::channel::mpsc;
use rust_decimal::Decimal;
use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use std::time::Duration;
#[cfg(feature = "native")]
use crate::account_syncer::{AccountSyncer, AccountSyncerConfig};
#[cfg(feature = "native")]
use crate::sync_traits::{AccountSync, TradeSync};
#[cfg(feature = "native")]
use crate::trade_syncer::{TradeSyncer, TradeSyncerConfig};
use crate::Engine;
use serde::Serialize;
const MAX_DEFERRED_ACTION_LIMIT_COMMANDS: usize = 1024;
const MAX_DEFERRED_ACTION_LIMIT_PER_LOOP: usize = 4;
#[derive(Debug, Clone)]
struct DeferredActionLimitCommand {
retry_at_ms: i64,
command: Command,
}
#[derive(Debug, Clone, Serialize)]
pub struct BacktestFill {
pub ts_ms: i64,
pub price: String,
pub qty: String,
pub side: String,
pub fee: String,
pub fee_amount: String,
pub fee_currency: String,
}
#[derive(Debug, Clone, Serialize)]
pub struct BacktestPositionSummary {
pub instrument: String,
pub final_position_qty: String,
pub avg_entry_price: Option<String>,
pub realized_pnl: String,
pub unrealized_pnl: Option<String>,
pub total_fees: String,
pub net_pnl: String,
}
#[derive(Debug, Clone, Serialize)]
pub struct BacktestResult {
pub fills: Vec<BacktestFill>,
pub trade_count: usize,
pub metrics: PerformanceMetrics,
pub benchmark: PerformanceBenchmark,
pub equity_curve: Vec<BacktestEquityPoint>,
pub closed_trades: Vec<BacktestClosedTrade>,
pub final_position_qty: String,
pub avg_entry_price: Option<String>,
pub realized_pnl: String,
pub unrealized_pnl: Option<String>,
pub total_fees: String,
pub total_volume: String,
pub net_pnl: String,
#[serde(default)]
pub positions: Vec<BacktestPositionSummary>,
pub exit_reason: Option<String>,
}
#[derive(Debug)]
pub enum PollResult {
Fills {
instance: ExchangeInstance,
fills: Vec<Fill>,
},
Quotes {
instance: ExchangeInstance,
quotes: Vec<Quote>,
},
ExchangeHealth {
instance: ExchangeInstance,
health: ExchangeHealth,
},
Error {
instance: ExchangeInstance,
error: String,
},
}
#[derive(Debug, Clone)]
pub struct RunnerConfig {
pub min_poll_delay_ms: u64,
pub initial_backoff_ms: u64,
pub max_backoff_ms: u64,
pub backoff_multiplier: f64,
pub quote_poll_interval_ms: u64,
pub cleanup_delay_ms: u64,
pub metrics_mode: String,
pub metrics_starting_balance_usdc: Option<Decimal>,
}
impl Default for RunnerConfig {
fn default() -> Self {
Self {
min_poll_delay_ms: 500,
initial_backoff_ms: 1000,
max_backoff_ms: 30_000,
backoff_multiplier: 2.0,
quote_poll_interval_ms: 1000,
cleanup_delay_ms: 5000, metrics_mode: "live".to_string(),
metrics_starting_balance_usdc: None,
}
}
}
#[derive(Debug, Default, Clone)]
pub struct TradingStats {
pub orders_placed: u32,
pub orders_filled: u32,
pub volume_traded: Decimal,
pub total_fees: Decimal,
pub realized_pnl: Decimal,
pub last_meta_log_ms: i64,
}
pub struct EngineRunner {
config: RunnerConfig,
engine: Engine,
exchanges: HashMap<ExchangeInstance, Arc<dyn Exchange>>,
instruments: Vec<InstrumentId>,
shutdown_rx: Option<mpsc::UnboundedReceiver<()>>,
shutdown_tx: mpsc::UnboundedSender<()>,
should_shutdown: bool,
shutdown_reason: Option<String>,
#[cfg(feature = "native")]
trade_syncer: Option<Box<dyn TradeSync>>,
#[cfg(feature = "native")]
account_syncer: Option<Box<dyn AccountSync>>,
current_mid_price: Option<Decimal>,
stats: TradingStats,
performance_tracker: PerformanceTracker,
fills_guard: PollGuard,
quotes_guard: PollGuard,
deferred_action_limit_commands: VecDeque<DeferredActionLimitCommand>,
}
impl EngineRunner {
pub fn new(engine: Engine, config: RunnerConfig) -> Self {
let (shutdown_tx, shutdown_rx) = mpsc::unbounded();
let fills_guard = PollGuard::new("fills", &config);
let quotes_guard = PollGuard::new("quotes", &config);
Self {
performance_tracker: PerformanceTracker::new(
config.metrics_mode.clone(),
config.metrics_starting_balance_usdc,
None,
),
config,
engine,
exchanges: HashMap::new(),
instruments: Vec::new(),
shutdown_rx: Some(shutdown_rx),
shutdown_tx,
should_shutdown: false,
shutdown_reason: None,
#[cfg(feature = "native")]
trade_syncer: None,
#[cfg(feature = "native")]
account_syncer: None,
current_mid_price: None,
stats: TradingStats::default(),
fills_guard,
quotes_guard,
deferred_action_limit_commands: VecDeque::new(),
}
}
#[cfg(feature = "native")]
pub fn with_syncer(mut self, syncer_config: TradeSyncerConfig) -> Self {
match TradeSyncer::new(syncer_config) {
Ok(syncer) => {
tracing::info!("[EngineRunner] Trade syncer enabled");
self.trade_syncer = Some(Box::new(syncer));
}
Err(e) => {
tracing::warn!("[EngineRunner] Failed to create trade syncer: {}", e);
}
}
self
}
#[cfg(feature = "native")]
pub fn set_trade_syncer(&mut self, syncer: Box<dyn TradeSync>) {
self.trade_syncer = Some(syncer);
}
#[cfg(feature = "native")]
pub fn with_account_syncer(mut self, syncer_config: AccountSyncerConfig) -> Self {
match AccountSyncer::new(syncer_config) {
Ok(syncer) => {
tracing::info!("[EngineRunner] Account syncer enabled");
self.account_syncer = Some(Box::new(syncer));
}
Err(e) => {
tracing::warn!("[EngineRunner] Failed to create account syncer: {}", e);
}
}
self
}
#[cfg(feature = "native")]
pub fn set_account_syncer(&mut self, syncer: Box<dyn AccountSync>) {
self.account_syncer = Some(syncer);
}
pub fn add_exchange(&mut self, exchange: Arc<dyn Exchange>) {
let instance = exchange.instance();
self.exchanges.insert(instance, exchange);
}
pub fn add_instrument(&mut self, instrument: InstrumentId) {
self.performance_tracker.set_instrument(instrument.clone());
self.instruments.push(instrument);
}
pub fn shutdown_handle(&self) -> mpsc::UnboundedSender<()> {
self.shutdown_tx.clone()
}
pub fn engine(&self) -> &Engine {
&self.engine
}
pub fn shutdown_reason(&self) -> Option<&str> {
self.shutdown_reason.as_deref()
}
pub fn get_backtest_results(&self, instrument: &InstrumentId) -> BacktestResult {
let fills = self.engine.get_fills();
let positions = self.tracked_positions();
let metrics_snapshot = self.performance_tracker.snapshot(fills, &positions);
let primary_position = self.engine.position(instrument);
let mut total_volume = Decimal::ZERO;
for fill in fills {
let notional = fill.qty.0 * fill.price.0;
total_volume += notional;
}
let total_realized_pnl: Decimal = positions.iter().map(|p| p.realized_pnl).sum();
let total_unrealized_pnl: Decimal = positions
.iter()
.map(|p| p.unrealized_pnl.unwrap_or_default())
.sum();
let total_fees: Decimal = positions.iter().map(|p| p.total_fees).sum();
let total_net_pnl: Decimal = positions.iter().map(Position::current_pnl).sum();
let unrealized_pnl = positions
.iter()
.any(|p| p.unrealized_pnl.is_some())
.then_some(total_unrealized_pnl);
let position_summaries: Vec<BacktestPositionSummary> = self
.instruments
.iter()
.zip(positions.iter())
.map(|(instrument, position)| BacktestPositionSummary {
instrument: instrument.to_string(),
final_position_qty: position.qty.to_string(),
avg_entry_price: position.avg_entry_px.map(|p| p.0.to_string()),
realized_pnl: position.realized_pnl.to_string(),
unrealized_pnl: position.unrealized_pnl.map(|p| p.to_string()),
total_fees: position.total_fees.to_string(),
net_pnl: position.current_pnl().to_string(),
})
.collect();
let backtest_fills: Vec<BacktestFill> = fills
.iter()
.map(|f| BacktestFill {
ts_ms: f.ts,
price: f.price.0.to_string(),
qty: f.qty.0.to_string(),
side: format!("{:?}", f.side),
fee: f.fee.asset.0.clone(),
fee_amount: f.fee.amount.to_string(),
fee_currency: f.fee.asset.0.clone(),
})
.collect();
BacktestResult {
trade_count: backtest_fills.len(),
fills: backtest_fills,
metrics: metrics_snapshot.metrics,
benchmark: metrics_snapshot.benchmark,
equity_curve: self.performance_tracker.equity_curve(),
closed_trades: self.performance_tracker.closed_trades(fills),
final_position_qty: primary_position.qty.to_string(),
avg_entry_price: primary_position.avg_entry_px.map(|p| p.0.to_string()),
realized_pnl: total_realized_pnl.to_string(),
unrealized_pnl: unrealized_pnl.map(|p| p.to_string()),
total_fees: total_fees.to_string(),
total_volume: total_volume.to_string(),
net_pnl: total_net_pnl.to_string(),
positions: position_summaries,
exit_reason: self.shutdown_reason.clone(),
}
}
fn tracked_positions(&self) -> Vec<Position> {
self.instruments
.iter()
.map(|instrument| self.engine.position(instrument))
.collect()
}
fn current_net_pnl(&self) -> Decimal {
self.tracked_positions()
.iter()
.map(Position::current_pnl)
.sum()
}
fn performance_snapshot(&self) -> PerformanceMetricsSnapshot {
let positions = self.tracked_positions();
self.performance_tracker
.snapshot(self.engine.get_fills(), &positions)
}
pub async fn run(&mut self) {
tracing::info!("Starting engine runner...");
for (instance, exchange) in &self.exchanges {
if let Err(e) = exchange.init().await {
tracing::error!("Exchange {} init failed: {}", instance, e);
self.shutdown_reason = Some(format!("exchange_init_failed:{}", e));
let stop_cmds = self.engine.stop_strategies();
self.execute_commands(stop_cmds).await;
return;
}
}
let start_cmds = self.engine.start_strategies();
self.execute_commands(start_cmds).await;
let mut shutdown_rx = self.shutdown_rx.take().expect("run called twice");
const MAX_EMPTY_POLLS: u32 = 3;
#[cfg(target_arch = "wasm32")]
let mut loop_iteration: u32 = 0;
loop {
#[cfg(target_arch = "wasm32")]
{
loop_iteration += 1;
}
#[cfg(target_arch = "wasm32")]
if loop_iteration <= 500 || loop_iteration % 10 == 0 {
web_sys::console::log_1(
&format!("[WASM Runner] Loop iteration {}", loop_iteration).into(),
);
}
if self.should_shutdown {
tracing::info!("Strategy requested shutdown");
break;
}
match shutdown_rx.try_next() {
Ok(Some(_)) | Ok(None) => {
tracing::info!("Shutdown signal received");
break;
}
Err(_) => {} }
self.process_deferred_action_limit_commands().await;
let exchanges_snapshot: Vec<(ExchangeInstance, Arc<dyn Exchange>)> = self
.exchanges
.iter()
.map(|(i, e)| (i.clone(), e.clone()))
.collect();
let sync_mechanism = self.engine.sync_mechanism();
for (instance, exchange) in &exchanges_snapshot {
match sync_mechanism {
bot_core::SyncMechanism::Poll => {
match self
.fills_guard
.execute(|| exchange.poll_user_fills(None))
.await
{
PollOutcome::Data(fills) => {
self.engine
.set_exchange_health(instance, ExchangeHealth::Active);
for fill in fills {
let client_id = self.resolve_fill_client_id(&fill);
if client_id.is_none() {
#[cfg(feature = "native")]
if let Some(ref mut syncer) = self.trade_syncer {
syncer.add_fill(fill.clone());
}
continue;
}
let client_id = client_id.expect("checked above");
if !self
.apply_fill_and_emit_events(instance, &client_id, &fill)
.await
{
continue;
}
#[cfg(feature = "native")]
if let Some(ref mut syncer) = self.trade_syncer {
syncer.add_fill(fill.clone());
}
tracing::info!(
"Fill: {} {} {} @ {} (oid={:?} tid={})",
fill.side,
fill.qty,
fill.instrument,
fill.price,
fill.exchange_order_id,
fill.trade_id
);
self.stats.orders_filled += 1;
self.stats.volume_traded += fill.qty.0 * fill.price.0;
self.stats.total_fees += fill.fee.amount;
}
}
PollOutcome::Empty => {} PollOutcome::Degraded(_) | PollOutcome::Fatal(_) => {
self.engine
.set_exchange_health(instance, ExchangeHealth::Halted);
}
}
}
bot_core::SyncMechanism::Snapshot => {
match self
.fills_guard
.execute(|| async {
exchange.poll_account_state().await.map(|state| vec![state])
})
.await
{
PollOutcome::Data(mut states) => {
let account_state = states.remove(0);
self.engine
.set_exchange_health(instance, ExchangeHealth::Active);
for pos_snapshot in &account_state.positions {
self.engine.apply_snapshot(
&pos_snapshot.instrument,
pos_snapshot.qty,
pos_snapshot.avg_entry_px,
pos_snapshot.unrealized_pnl,
);
tracing::debug!(
"Position snapshot: {} qty={} entry={:?} pnl={:?}",
pos_snapshot.instrument,
pos_snapshot.qty,
pos_snapshot.avg_entry_px,
pos_snapshot.unrealized_pnl
);
}
#[cfg(feature = "native")]
let account_metrics_snapshot = if self.account_syncer.is_some() {
Some(self.performance_snapshot())
} else {
None
};
#[cfg(feature = "native")]
if let Some(ref mut syncer) = self.account_syncer {
syncer.set_metrics_snapshot(account_metrics_snapshot);
if syncer.should_sync() {
if let Err(e) = syncer.sync(&account_state, false, "").await
{
tracing::error!(
"[{}] Account sync failed: {}",
instance,
e
);
}
}
}
tracing::info!(
"Account snapshot: positions={} account_value={:?} pnl={:?}",
account_state.positions.len(),
account_state.account_value,
account_state.unrealized_pnl
);
}
PollOutcome::Empty => {} PollOutcome::Degraded(_) | PollOutcome::Fatal(_) => {
self.engine
.set_exchange_health(instance, ExchangeHealth::Halted);
}
}
}
}
}
if !self.instruments.is_empty() {
for (instance, exchange) in &exchanges_snapshot {
match self
.quotes_guard
.execute(|| exchange.poll_quotes(&self.instruments))
.await
{
PollOutcome::Data(quotes) => {
self.engine
.set_exchange_health(instance, ExchangeHealth::Active);
for quote in quotes {
self.engine.update_quote(quote.clone());
let mid = (quote.bid.0 + quote.ask.0) / Decimal::TWO;
self.current_mid_price = Some(mid);
tracing::debug!(
"[EngineRunner] Updated mid price: {} (bid={}, ask={})",
mid,
quote.bid,
quote.ask
);
let event = Event::Quote(QuoteEvent {
exchange: instance.exchange_id.clone(),
instrument: quote.instrument.clone(),
bid: quote.bid,
ask: quote.ask,
ts: quote.ts,
});
self.handle_event(event).await;
self.performance_tracker.record_equity_point(
quote.ts,
Some(mid),
self.current_net_pnl(),
);
}
}
PollOutcome::Empty => {} PollOutcome::Degraded(_) | PollOutcome::Fatal(_) => {
self.engine
.set_exchange_health(instance, ExchangeHealth::Halted);
}
}
}
}
if self.config.min_poll_delay_ms == 0
&& self.quotes_guard.looks_exhausted(MAX_EMPTY_POLLS)
{
tracing::info!("[EngineRunner] Backtest complete: no more quotes available");
#[cfg(target_arch = "wasm32")]
web_sys::console::log_1(&"[WASM Runner] Backtest complete - exiting".into());
self.shutdown_reason = Some("Backtest complete - all quotes processed".to_string());
break;
}
#[cfg(feature = "native")]
let trade_metrics_snapshot = if self.trade_syncer.is_some() {
Some(self.performance_snapshot())
} else {
None
};
#[cfg(feature = "native")]
if let Some(ref mut syncer) = self.trade_syncer {
syncer.set_metrics_snapshot(trade_metrics_snapshot);
if syncer.should_sync() {
match syncer.sync(self.current_mid_price, false, "").await {
Ok(result) => {
if let Some(pnl) = result.pnl {
tracing::debug!("[EngineRunner] Sync success: pnl={:.4}", pnl);
}
}
Err(e) => {
tracing::warn!("[EngineRunner] Sync failed: {}", e);
}
}
}
}
let now_ms = bot_core::now_ms();
if now_ms - self.stats.last_meta_log_ms >= 30_000 {
self.stats.last_meta_log_ms = now_ms;
let positions: Vec<(InstrumentId, Position)> = self
.instruments
.iter()
.map(|i| (i.clone(), self.engine.position(i)))
.collect();
let pos_str = positions
.iter()
.map(|(instrument, position)| format!("{}:{:.4}", instrument, position.qty))
.collect::<Vec<_>>()
.join("/");
let pos_detail_str = positions
.iter()
.map(|(instrument, position)| {
let avg = position
.avg_entry_px
.map(|price| price.0.to_string())
.unwrap_or_else(|| "-".to_string());
let unrealized = position.unrealized_pnl.unwrap_or_default();
format!(
"{}[qty={:.4},avg={},r_pnl={:.4},u_pnl={:.4},net={:.4}]",
instrument,
position.qty,
avg,
position.realized_pnl,
unrealized,
position.current_pnl()
)
})
.collect::<Vec<_>>()
.join(" ");
let total_realized_pnl: Decimal =
positions.iter().map(|(_, p)| p.realized_pnl).sum();
let total_unrealized_pnl: Decimal = positions
.iter()
.map(|(_, p)| p.unrealized_pnl.unwrap_or_default())
.sum();
let total_position_fees: Decimal =
positions.iter().map(|(_, p)| p.total_fees).sum();
let total_net_pnl: Decimal = positions.iter().map(|(_, p)| p.current_pnl()).sum();
let sync_mechanism = self.engine.sync_mechanism();
if sync_mechanism == bot_core::SyncMechanism::Poll {
tracing::info!(
"[META] pos={} orders_placed={} orders_filled={} vol={:.2} fees={:.4} r_pnl={:.4} u_pnl={:.4} net_pnl={:.4} legs={}",
pos_str,
self.stats.orders_placed,
self.stats.orders_filled,
self.stats.volume_traded,
self.stats.total_fees,
total_realized_pnl,
total_unrealized_pnl,
total_net_pnl,
pos_detail_str
);
} else {
tracing::info!(
"[META] pos={} r_pnl={:.4} u_pnl={:.4} fees={:.4} net_pnl={:.4} legs={} (snapshot)",
pos_str,
total_realized_pnl,
total_unrealized_pnl,
total_position_fees,
total_net_pnl,
pos_detail_str
);
}
}
if self.config.min_poll_delay_ms > 0 {
sleep(Duration::from_millis(self.config.min_poll_delay_ms)).await;
} else {
crate::compat::yield_now().await;
}
}
let stop_cmds = self.engine.stop_strategies();
self.execute_commands(stop_cmds).await;
if self.config.cleanup_delay_ms > 0 {
tracing::info!(
"Waiting {}ms for cleanup to complete...",
self.config.cleanup_delay_ms
);
sleep(Duration::from_millis(self.config.cleanup_delay_ms)).await;
}
#[cfg(feature = "native")]
let final_trade_metrics_snapshot = if self.trade_syncer.is_some() {
Some(self.performance_snapshot())
} else {
None
};
#[cfg(feature = "native")]
if let Some(ref mut syncer) = self.trade_syncer {
syncer.set_metrics_snapshot(final_trade_metrics_snapshot);
let reason = self
.shutdown_reason
.as_deref()
.unwrap_or("shutdown:graceful");
tracing::info!(
"[EngineRunner] Performing final sync before shutdown with reason: {}",
reason
);
match syncer.shutdown_sync(self.current_mid_price, reason).await {
Ok(result) => {
tracing::info!("[EngineRunner] Final sync complete: pnl={:?}", result.pnl);
}
Err(e) => {
tracing::warn!("[EngineRunner] Final sync failed: {}", e);
}
}
}
if let Some(first_instrument) = self.instruments.first() {
let results = self.get_backtest_results(first_instrument);
if let Ok(json) = serde_json::to_string(&results) {
println!("{}", json);
use std::io::Write;
let _ = std::io::stdout().flush();
}
}
tracing::info!("Engine runner stopped");
}
async fn handle_event(&mut self, event: Event) {
use std::collections::VecDeque;
let mut queue: VecDeque<Event> = VecDeque::new();
queue.push_back(event);
while let Some(ev) = queue.pop_front() {
let cmds = self.engine.dispatch_event(&ev);
let followups = self.execute_commands(cmds).await;
for f in followups {
queue.push_back(f);
}
}
}
fn action_limit_retry_after(error: &ExchangeError) -> Option<u64> {
match error {
ExchangeError::WouldExceedUserActionLimit { retry_after_ms, .. } => {
Some(*retry_after_ms)
}
_ => None,
}
}
fn is_cancel_command(command: &Command) -> bool {
matches!(command, Command::CancelOrder(_) | Command::CancelAll(_))
}
fn same_cancel_command(left: &Command, right: &Command) -> bool {
match (left, right) {
(Command::CancelOrder(a), Command::CancelOrder(b)) => {
a.exchange == b.exchange && a.client_id == b.client_id
}
(Command::CancelAll(a), Command::CancelAll(b)) => {
a.exchange == b.exchange && a.instrument == b.instrument
}
_ => false,
}
}
fn defer_action_limit_command(&mut self, command: Command, retry_after_ms: u64) {
let retry_at_ms = bot_core::now_ms() + retry_after_ms as i64;
if Self::is_cancel_command(&command) {
if let Some(existing) = self
.deferred_action_limit_commands
.iter_mut()
.find(|existing| Self::same_cancel_command(&existing.command, &command))
{
existing.retry_at_ms = existing.retry_at_ms.min(retry_at_ms);
tracing::debug!(
"Coalesced duplicate Hyperliquid action-limit cancel command: {:?}",
command
);
return;
}
}
if self.deferred_action_limit_commands.len() >= MAX_DEFERRED_ACTION_LIMIT_COMMANDS {
tracing::error!(
"Action-limit deferred queue full ({} commands); stopping runner",
self.deferred_action_limit_commands.len()
);
self.shutdown_reason = Some("action_limit_deferred_queue_full".to_string());
self.should_shutdown = true;
return;
}
tracing::warn!(
"Deferring command for Hyperliquid action limit: retry_after_ms={} command={:?}",
retry_after_ms,
command
);
self.deferred_action_limit_commands
.push_back(DeferredActionLimitCommand {
retry_at_ms,
command,
});
}
fn next_deferred_action_limit_index(&self, now: i64) -> Option<usize> {
let mut selected: Option<(usize, bool, i64)> = None;
for (index, deferred) in self.deferred_action_limit_commands.iter().enumerate() {
if deferred.retry_at_ms > now {
continue;
}
let is_cancel = Self::is_cancel_command(&deferred.command);
match selected {
None => selected = Some((index, is_cancel, deferred.retry_at_ms)),
Some((_, selected_is_cancel, selected_retry_at_ms))
if (is_cancel && !selected_is_cancel)
|| (is_cancel == selected_is_cancel
&& deferred.retry_at_ms < selected_retry_at_ms) =>
{
selected = Some((index, is_cancel, deferred.retry_at_ms));
}
_ => {}
}
}
selected.map(|(index, _, _)| index)
}
async fn process_deferred_action_limit_commands(&mut self) {
let now = bot_core::now_ms();
let mut processed = 0usize;
while processed < MAX_DEFERRED_ACTION_LIMIT_PER_LOOP {
let Some(index) = self.next_deferred_action_limit_index(now) else {
break;
};
let deferred = self
.deferred_action_limit_commands
.remove(index)
.expect("index selected from queue");
tracing::info!(
"Retrying Hyperliquid action-limit deferred command: {:?}",
deferred.command
);
let followups = self.execute_commands(vec![deferred.command]).await;
for event in followups {
self.handle_event(event).await;
}
processed += 1;
}
}
fn resolve_fill_client_id(&self, fill: &Fill) -> Option<ClientOrderId> {
if let Some(cid) = fill.client_id.clone() {
return Some(cid);
}
if let Some(ref eid) = fill.exchange_order_id {
if let Some(cid) = self.engine.order_manager().client_id_from_exchange_id(eid) {
return Some(cid.clone());
}
}
None
}
async fn apply_fill_and_emit_events(
&mut self,
instance: &ExchangeInstance,
client_id: &ClientOrderId,
fill: &Fill,
) -> bool {
let is_new = self.engine.order_manager_mut().apply_fill(
client_id,
&fill.trade_id,
fill.qty,
fill.price,
);
if !is_new {
return false;
}
let net_qty = if let Some(meta) = self.engine.instrument_meta(&fill.instrument) {
if meta.kind == bot_core::InstrumentKind::Spot
&& fill.side == bot_core::OrderSide::Buy
&& fill.fee.asset == meta.base_asset
{
let nq = Qty::new((fill.qty.0 - fill.fee.amount).max(Decimal::ZERO));
tracing::debug!(
"Spot BUY fee deduction: gross={} fee={} net={}",
fill.qty,
fill.fee.amount,
nq
);
nq
} else {
fill.qty
}
} else {
fill.qty
};
self.engine
.apply_position_fill(&fill.instrument, fill.side, net_qty, fill.price);
self.engine
.apply_fill_fee(&fill.instrument, fill.fee.amount);
let filled_event = Event::OrderFilled(OrderFilledEvent {
exchange: instance.exchange_id.clone(),
instrument: fill.instrument.clone(),
client_id: client_id.clone(),
trade_id: fill.trade_id.clone(),
side: fill.side,
price: fill.price,
qty: fill.qty,
net_qty,
fee: fill.fee.clone(),
ts: fill.ts,
});
if let Event::OrderFilled(ref fill_event) = filled_event {
self.engine.record_fill(fill_event.clone());
}
self.handle_event(filled_event).await;
let is_complete = self.engine.order_manager().is_complete(client_id);
if is_complete {
if let Some(order) = self.engine.order_manager().get(client_id).cloned() {
let completed_event = Event::OrderCompleted(OrderCompletedEvent {
exchange: instance.exchange_id.clone(),
instrument: order.instrument.clone(),
client_id: client_id.clone(),
filled_qty: order.filled_qty,
avg_fill_px: order.avg_fill_px,
ts: bot_core::now_ms(),
});
self.handle_event(completed_event).await;
}
}
true
}
async fn execute_commands(&mut self, cmds: Vec<Command>) -> Vec<Event> {
let mut followups: Vec<Event> = Vec::new();
for cmd in cmds {
match cmd {
Command::PlaceOrder(c) => {
let mut evs = self.execute_place(c).await;
followups.append(&mut evs);
}
Command::PlaceOrders(orders) => {
let mut evs = self.execute_place_batch(orders).await;
followups.append(&mut evs);
}
Command::CancelOrder(c) => {
if let Some(ev) = self.execute_cancel(c).await {
followups.push(ev);
}
}
Command::CancelAll(c) => {
let mut evs = self.execute_cancel_all(c).await;
followups.append(&mut evs);
}
Command::StopStrategy(stop) => {
tracing::info!("Strategy {} stopped: {}", stop.strategy_id.0, stop.reason);
self.shutdown_reason = Some(stop.reason.clone());
self.should_shutdown = true;
}
}
}
followups
}
async fn execute_place(&mut self, cmd: PlaceOrder) -> Vec<Event> {
#[cfg(feature = "wasm")]
web_sys::console::log_1(
&format!(
"[Runner] execute_place called: instrument={}, exchange={:?}",
cmd.instrument, cmd.exchange
)
.into(),
);
let Some(exchange) = self.exchanges.get(&cmd.exchange).cloned() else {
#[cfg(feature = "wasm")]
{
let available: Vec<_> = self.exchanges.keys().collect();
web_sys::console::log_1(
&format!(
"[Runner] Exchange NOT FOUND: {:?}, available: {:?}",
cmd.exchange, available
)
.into(),
);
}
return vec![];
};
#[cfg(feature = "wasm")]
web_sys::console::log_1(&format!("[Runner] Exchange found, proceeding with order").into());
let (market_index, rounded_price, rounded_qty) =
match self.engine.instrument_meta(&cmd.instrument) {
Some(meta) => {
let rp = meta.round_price(cmd.price);
let rq = meta.round_qty(cmd.qty);
tracing::debug!(
"Rounding: price {} -> {}, qty {} -> {}",
cmd.price,
rp,
cmd.qty,
rq
);
(meta.market_index.clone(), rp, rq)
}
None => {
tracing::warn!(
"PlaceOrder ignored: no InstrumentMeta for {}",
cmd.instrument
);
return vec![Event::OrderRejected(OrderRejectedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
reason: "missing instrument meta".to_string(),
ts: bot_core::now_ms(),
})];
}
};
self.engine.order_manager_mut().create_order(
cmd.client_id.clone(),
cmd.instrument.clone(),
cmd.side,
rounded_price,
rounded_qty,
);
let order_input = OrderInput {
instrument: cmd.instrument.clone(),
market_index,
client_id: cmd.client_id.clone(),
side: cmd.side,
price: rounded_price,
qty: rounded_qty,
tif: cmd.tif,
post_only: cmd.post_only,
reduce_only: cmd.reduce_only,
};
match exchange.place_orders(&[order_input]).await {
Ok(results) => {
match results.into_iter().next() {
Some(PlaceOrderResult::Accepted {
exchange_order_id,
filled_qty,
avg_fill_px,
}) => {
self.engine
.order_manager_mut()
.accept_order(&cmd.client_id, exchange_order_id.clone());
self.stats.orders_placed += 1;
let accepted_event = Event::OrderAccepted(OrderAcceptedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
exchange_order_id,
ts: bot_core::now_ms(),
});
let use_ioc_fill =
self.engine.sync_mechanism() == bot_core::SyncMechanism::Snapshot;
if use_ioc_fill {
if let (Some(qty), Some(px)) = (&filled_qty, &avg_fill_px) {
vec![
accepted_event,
Event::OrderCompleted(OrderCompletedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
filled_qty: *qty,
avg_fill_px: Some(*px),
ts: bot_core::now_ms(),
}),
]
} else {
vec![accepted_event]
}
} else {
vec![accepted_event]
}
}
Some(PlaceOrderResult::Rejected { reason }) => {
self.engine.order_manager_mut().reject_order(&cmd.client_id);
vec![Event::OrderRejected(OrderRejectedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
reason,
ts: bot_core::now_ms(),
})]
}
None => {
self.engine.order_manager_mut().reject_order(&cmd.client_id);
vec![Event::OrderRejected(OrderRejectedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
reason: "No result returned from place_orders".to_string(),
ts: bot_core::now_ms(),
})]
}
}
}
Err(e) => {
if let Some(retry_after_ms) = Self::action_limit_retry_after(&e) {
self.defer_action_limit_command(Command::PlaceOrder(cmd), retry_after_ms);
return vec![];
}
self.engine.order_manager_mut().reject_order(&cmd.client_id);
vec![Event::OrderRejected(OrderRejectedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
reason: e.to_string(),
ts: bot_core::now_ms(),
})]
}
}
}
async fn execute_place_batch(&mut self, orders: Vec<PlaceOrder>) -> Vec<Event> {
if orders.is_empty() {
return Vec::new();
}
let exchange_instance = &orders[0].exchange;
let exchange = match self.exchanges.get(exchange_instance) {
Some(e) => e.clone(),
None => {
return orders
.iter()
.map(|cmd| {
Event::OrderRejected(OrderRejectedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
reason: "Exchange not found".to_string(),
ts: bot_core::now_ms(),
})
})
.collect();
}
};
let mut order_inputs: Vec<OrderInput> = Vec::with_capacity(orders.len());
let mut cmd_metadata: Vec<(PlaceOrder, bool)> = Vec::with_capacity(orders.len());
for cmd in orders {
let meta_result = self.engine.instrument_meta(&cmd.instrument).map(|meta| {
let rp = meta.round_price(cmd.price);
let rq = meta.round_qty(cmd.qty);
(meta.market_index.clone(), rp, rq)
});
match meta_result {
Some((market_index, rounded_price, rounded_qty)) => {
self.engine.order_manager_mut().create_order(
cmd.client_id.clone(),
cmd.instrument.clone(),
cmd.side,
rounded_price,
rounded_qty,
);
order_inputs.push(OrderInput {
instrument: cmd.instrument.clone(),
market_index,
client_id: cmd.client_id.clone(),
side: cmd.side,
price: rounded_price,
qty: rounded_qty,
tif: cmd.tif,
post_only: cmd.post_only,
reduce_only: cmd.reduce_only,
});
cmd_metadata.push((cmd, true));
}
None => {
tracing::warn!(
"PlaceOrder ignored: no InstrumentMeta for {}",
cmd.instrument
);
cmd_metadata.push((cmd, false));
}
}
}
let mut events: Vec<Event> = cmd_metadata
.iter()
.filter(|(_, valid)| !valid)
.map(|(cmd, _)| {
Event::OrderRejected(OrderRejectedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
reason: "missing instrument meta".to_string(),
ts: bot_core::now_ms(),
})
})
.collect();
if order_inputs.is_empty() {
return events;
}
tracing::info!(
"Executing batch place_orders with {} orders",
order_inputs.len()
);
match exchange.place_orders(&order_inputs).await {
Ok(results) => {
let valid_cmds: Vec<&PlaceOrder> = cmd_metadata
.iter()
.filter(|(_, valid)| *valid)
.map(|(cmd, _)| cmd)
.collect();
for (i, result) in results.into_iter().enumerate() {
let cmd = valid_cmds.get(i);
if let Some(cmd) = cmd {
match result {
PlaceOrderResult::Accepted {
exchange_order_id,
filled_qty,
avg_fill_px,
} => {
self.engine
.order_manager_mut()
.accept_order(&cmd.client_id, exchange_order_id.clone());
let use_ioc_fill = self.engine.sync_mechanism()
== bot_core::SyncMechanism::Snapshot;
if use_ioc_fill {
if let (Some(qty), Some(px)) = (&filled_qty, &avg_fill_px) {
events.push(Event::OrderCompleted(OrderCompletedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
filled_qty: *qty,
avg_fill_px: Some(*px),
ts: bot_core::now_ms(),
}));
} else {
events.push(Event::OrderAccepted(OrderAcceptedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
exchange_order_id,
ts: bot_core::now_ms(),
}));
}
} else {
events.push(Event::OrderAccepted(OrderAcceptedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
exchange_order_id,
ts: bot_core::now_ms(),
}));
}
}
PlaceOrderResult::Rejected { reason } => {
self.engine.order_manager_mut().reject_order(&cmd.client_id);
events.push(Event::OrderRejected(OrderRejectedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
reason,
ts: bot_core::now_ms(),
}));
}
}
}
}
}
Err(e) => {
if let Some(retry_after_ms) = Self::action_limit_retry_after(&e) {
let deferred_orders: Vec<PlaceOrder> = cmd_metadata
.iter()
.filter(|(_, valid)| *valid)
.map(|(cmd, _)| cmd.clone())
.collect();
if !deferred_orders.is_empty() {
self.defer_action_limit_command(
Command::PlaceOrders(deferred_orders),
retry_after_ms,
);
}
return events;
}
for (cmd, valid) in &cmd_metadata {
if *valid {
self.engine.order_manager_mut().reject_order(&cmd.client_id);
events.push(Event::OrderRejected(OrderRejectedEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument: cmd.instrument.clone(),
client_id: cmd.client_id.clone(),
reason: e.to_string(),
ts: bot_core::now_ms(),
}));
}
}
}
}
events
}
async fn execute_cancel(&mut self, cmd: CancelOrder) -> Option<Event> {
let exchange = self.exchanges.get(&cmd.exchange)?.clone();
let (instrument, exchange_order_id) = match self.engine.order_manager().get(&cmd.client_id)
{
Some(o) => (o.instrument.clone(), o.exchange_order_id.clone()),
None => {
tracing::warn!("CancelOrder ignored: unknown client_id {}", cmd.client_id);
return None;
}
};
let market_index = match self.engine.instrument_meta(&instrument) {
Some(m) => m.market_index.clone(),
None => {
tracing::warn!("CancelOrder ignored: no InstrumentMeta for {}", instrument);
return None;
}
};
match exchange
.cancel_order(
&instrument,
&market_index,
&cmd.client_id,
exchange_order_id.as_ref(),
)
.await
{
Ok(_) => {
self.engine.order_manager_mut().cancel_order(&cmd.client_id);
Some(Event::OrderCanceled(OrderCanceledEvent {
exchange: cmd.exchange.exchange_id.clone(),
instrument,
client_id: cmd.client_id.clone(),
reason: None,
ts: bot_core::now_ms(),
}))
}
Err(e) => {
if let Some(retry_after_ms) = Self::action_limit_retry_after(&e) {
self.defer_action_limit_command(Command::CancelOrder(cmd), retry_after_ms);
return None;
}
tracing::warn!("CancelOrder failed for {}: {}", cmd.client_id, e);
None
}
}
}
async fn execute_cancel_all(&mut self, cmd: CancelAll) -> Vec<Event> {
tracing::info!(
"Executing CancelAll for exchange={:?} instrument={:?}",
cmd.exchange.exchange_id,
cmd.instrument.as_ref().map(|i| i.to_string())
);
let out = Vec::new();
let exchange = match self.exchanges.get(&cmd.exchange) {
Some(e) => e.clone(),
None => {
tracing::warn!(
"CancelAll: exchange not found {:?}",
cmd.exchange.exchange_id
);
return out;
}
};
let instruments: Vec<InstrumentId> = if let Some(i) = cmd.instrument.clone() {
vec![i]
} else {
self.instruments.clone()
};
for instrument in instruments {
let market_index = match self.engine.instrument_meta(&instrument) {
Some(m) => m.market_index.clone(),
None => continue,
};
tracing::info!("CancelAll: calling cancel_all_orders for {}", instrument);
match exchange.cancel_all_orders(&instrument, &market_index).await {
Ok(n) => {
tracing::info!(
"CancelAll: successfully canceled {} orders for {}",
n,
instrument
);
}
Err(e) => {
if let Some(retry_after_ms) = Self::action_limit_retry_after(&e) {
self.defer_action_limit_command(
Command::CancelAll(cmd.clone()),
retry_after_ms,
);
return out;
}
tracing::warn!("CancelAll failed for {}: {}", instrument, e);
}
}
}
out
}
}
#[cfg(feature = "native")]
pub fn spawn_runner(
engine: Engine,
exchanges: Vec<Arc<dyn Exchange>>,
instruments: Vec<InstrumentId>,
config: RunnerConfig,
) -> (tokio::task::JoinHandle<()>, mpsc::UnboundedSender<()>) {
spawn_runner_with_syncer(engine, exchanges, instruments, config, None, None)
}
#[cfg(feature = "native")]
pub fn spawn_runner_with_syncer(
engine: Engine,
exchanges: Vec<Arc<dyn Exchange>>,
instruments: Vec<InstrumentId>,
config: RunnerConfig,
syncer_config: Option<TradeSyncerConfig>,
account_syncer_config: Option<AccountSyncerConfig>,
) -> (tokio::task::JoinHandle<()>, mpsc::UnboundedSender<()>) {
let mut runner = EngineRunner::new(engine, config);
if let Some(syncer_cfg) = syncer_config {
runner = runner.with_syncer(syncer_cfg);
}
if let Some(account_cfg) = account_syncer_config {
runner = runner.with_account_syncer(account_cfg);
}
for exchange in exchanges {
runner.add_exchange(exchange);
}
for instrument in instruments {
runner.add_instrument(instrument);
}
let shutdown_handle = runner.shutdown_handle();
let handle = tokio::spawn(async move {
runner.run().await;
});
(handle, shutdown_handle)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::testing::MockExchange;
use crate::EngineConfig;
use bot_core::{
AssetId, Environment, ExchangeId, Fee, InstrumentKind, InstrumentMeta, MarketIndex,
OrderSide, Price, Quote, Strategy, StrategyContext, StrategyId, TimeInForce, TradeId,
};
use rust_decimal::Decimal;
use rust_decimal_macros::dec;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::time::timeout;
struct DeferredPlacementProbeStrategy {
id: StrategyId,
exchange: ExchangeInstance,
instrument: InstrumentId,
accepted: Arc<AtomicUsize>,
rejected: Arc<AtomicUsize>,
single_sent: bool,
}
impl DeferredPlacementProbeStrategy {
fn new(
exchange: ExchangeInstance,
instrument: InstrumentId,
accepted: Arc<AtomicUsize>,
rejected: Arc<AtomicUsize>,
) -> Self {
Self {
id: StrategyId::new("deferred-placement-probe"),
exchange,
instrument,
accepted,
rejected,
single_sent: false,
}
}
fn order(&self, client_id: &str) -> PlaceOrder {
PlaceOrder {
client_id: ClientOrderId::new(client_id),
exchange: self.exchange.clone(),
instrument: self.instrument.clone(),
side: OrderSide::Buy,
price: Price::new(Decimal::new(100, 0)),
qty: Qty::new(Decimal::new(1, 0)),
tif: TimeInForce::Gtc,
post_only: false,
reduce_only: false,
}
}
}
impl Strategy for DeferredPlacementProbeStrategy {
fn id(&self) -> &StrategyId {
&self.id
}
fn on_start(&mut self, ctx: &mut dyn StrategyContext) {
ctx.place_orders(vec![self.order("batch-1"), self.order("batch-2")]);
}
fn on_event(&mut self, ctx: &mut dyn StrategyContext, event: &Event) {
match event {
Event::OrderAccepted(_) => {
let accepted = self.accepted.fetch_add(1, Ordering::SeqCst) + 1;
if accepted == 2 && !self.single_sent {
self.single_sent = true;
ctx.place_order(self.order("later-single"));
}
if accepted >= 3 {
ctx.stop_strategy(self.id.clone(), "probe complete");
}
}
Event::OrderRejected(_) => {
self.rejected.fetch_add(1, Ordering::SeqCst);
}
_ => {}
}
}
fn on_timer(&mut self, _ctx: &mut dyn StrategyContext, _timer_id: bot_core::TimerId) {}
fn on_stop(&mut self, _ctx: &mut dyn StrategyContext) {}
}
fn action_limit_error(needed: u32) -> ExchangeError {
ExchangeError::WouldExceedUserActionLimit {
retry_after_ms: 1,
needed,
}
}
fn instrument_meta(instrument: &InstrumentId) -> InstrumentMeta {
InstrumentMeta {
instrument_id: instrument.clone(),
market_index: MarketIndex::new(0),
base_asset: AssetId::new("BTC"),
quote_asset: AssetId::new("USDC"),
tick_size: Decimal::new(1, 0),
lot_size: Decimal::new(1, 0),
min_qty: None,
min_notional: None,
fee_asset_default: Some(AssetId::new("USDC")),
kind: InstrumentKind::Perp,
}
}
fn record_backtest_fill(
runner: &mut EngineRunner,
instrument: &InstrumentId,
side: OrderSide,
price: Decimal,
qty: Decimal,
fee: Decimal,
ts: i64,
trade_id: &str,
) {
runner
.engine
.apply_position_fill(instrument, side, Qty::new(qty), Price::new(price));
runner.engine.apply_fill_fee(instrument, fee);
runner.engine.record_fill(OrderFilledEvent {
exchange: ExchangeId::new("hyperliquid"),
trade_id: TradeId::new(trade_id),
client_id: ClientOrderId::new(format!("client-{}", trade_id)),
instrument: instrument.clone(),
side,
price: Price::new(price),
qty: Qty::new(qty),
net_qty: Qty::new(qty),
fee: Fee::new(fee, AssetId::new("USDC")),
ts,
});
}
fn record_backtest_equity(
runner: &mut EngineRunner,
instrument: &InstrumentId,
price: Decimal,
ts: i64,
) {
runner.engine.update_quote(Quote {
instrument: instrument.clone(),
bid: Price::new(price),
ask: Price::new(price),
bid_size: Qty::new(dec!(1000)),
ask_size: Qty::new(dec!(1000)),
ts,
});
runner
.performance_tracker
.record_equity_point(ts, Some(price), runner.current_net_pnl());
}
#[test]
fn backtest_result_embeds_deterministic_performance_metrics() {
let instrument = InstrumentId::new("BTC-PERP");
let mut engine = Engine::new(EngineConfig::default());
engine.register_instrument(instrument_meta(&instrument));
let mut runner = EngineRunner::new(
engine,
RunnerConfig {
metrics_mode: "backtest".to_string(),
metrics_starting_balance_usdc: Some(dec!(1000)),
cleanup_delay_ms: 0,
..Default::default()
},
);
runner.add_instrument(instrument.clone());
record_backtest_equity(&mut runner, &instrument, dec!(100), 1);
record_backtest_fill(
&mut runner,
&instrument,
OrderSide::Buy,
dec!(100),
dec!(1),
dec!(1),
2,
"open-long",
);
record_backtest_equity(&mut runner, &instrument, dec!(100), 2);
record_backtest_fill(
&mut runner,
&instrument,
OrderSide::Sell,
dec!(110),
dec!(1),
dec!(1),
3,
"close-long",
);
record_backtest_equity(&mut runner, &instrument, dec!(110), 3);
let result = runner.get_backtest_results(&instrument);
assert_eq!(result.trade_count, 2);
assert_eq!(result.metrics.fill_count, 2);
assert_eq!(result.metrics.closed_trade_count, 1);
assert_eq!(result.metrics.win_rate_pct, Some(100.0));
assert_eq!(result.metrics.winning_trade_count, 1);
assert_eq!(result.metrics.losing_trade_count, 0);
assert_eq!(result.metrics.total_fees, "2");
assert_eq!(result.metrics.total_volume, "210");
assert_eq!(result.metrics.net_pnl, "8");
assert_eq!(result.realized_pnl, "10");
assert_eq!(result.net_pnl, "8");
assert_eq!(result.positions.len(), 1);
assert_eq!(result.positions[0].instrument, "BTC-PERP");
assert_eq!(result.positions[0].net_pnl, "8");
assert_eq!(result.metrics.period_return_pct, Some(0.8));
assert_eq!(result.metrics.max_drawdown_usdc, "1");
assert_eq!(result.metrics.max_drawdown_pct, Some(0.1));
assert_eq!(result.fills[0].fee, "USDC");
assert_eq!(result.fills[0].fee_amount, "1");
assert_eq!(result.fills[0].fee_currency, "USDC");
assert_eq!(result.closed_trades.len(), 1);
assert_eq!(result.closed_trades[0].gross_pnl, "10");
assert_eq!(result.closed_trades[0].fees, "2");
assert_eq!(result.closed_trades[0].net_pnl, "8");
assert_eq!(result.equity_curve.len(), 3);
assert_eq!(result.benchmark.quote_count, 3);
assert_eq!(
result.benchmark.starting_balance_usdc,
Some("1000".to_string())
);
assert_eq!(
result.benchmark.ending_balance_usdc,
Some("1008".to_string())
);
assert_eq!(result.benchmark.instrument, Some("BTC-PERP".to_string()));
}
#[test]
fn backtest_result_aggregates_multi_instrument_pnl() {
let yes = InstrumentId::new("#50-OUTCOME");
let no = InstrumentId::new("#51-OUTCOME");
let mut engine = Engine::new(EngineConfig::default());
engine.register_instrument(instrument_meta(&yes));
engine.register_instrument(instrument_meta(&no));
let mut runner = EngineRunner::new(
engine,
RunnerConfig {
metrics_mode: "backtest".to_string(),
metrics_starting_balance_usdc: Some(dec!(100)),
cleanup_delay_ms: 0,
..Default::default()
},
);
runner.add_instrument(yes.clone());
runner.add_instrument(no.clone());
record_backtest_fill(
&mut runner,
&yes,
OrderSide::Buy,
dec!(1),
dec!(10),
Decimal::ZERO,
1,
"yes-open",
);
record_backtest_fill(
&mut runner,
&yes,
OrderSide::Sell,
dec!(2),
dec!(5),
Decimal::ZERO,
2,
"yes-partial-close",
);
record_backtest_fill(
&mut runner,
&no,
OrderSide::Buy,
dec!(4),
dec!(10),
Decimal::ZERO,
3,
"no-open",
);
record_backtest_fill(
&mut runner,
&no,
OrderSide::Sell,
dec!(5),
dec!(10),
Decimal::ZERO,
4,
"no-close",
);
record_backtest_equity(&mut runner, &yes, dec!(3), 5);
record_backtest_equity(&mut runner, &no, dec!(5), 5);
let result = runner.get_backtest_results(&yes);
assert_eq!(result.realized_pnl, "15");
assert_eq!(result.unrealized_pnl, Some("10".to_string()));
assert_eq!(result.net_pnl, "25");
assert_eq!(result.positions.len(), 2);
let yes_position = result
.positions
.iter()
.find(|position| position.instrument == "#50-OUTCOME")
.expect("yes position summary");
assert_eq!(yes_position.final_position_qty, "5");
assert_eq!(yes_position.realized_pnl, "5");
assert_eq!(yes_position.unrealized_pnl, Some("10".to_string()));
assert_eq!(yes_position.net_pnl, "15");
let no_position = result
.positions
.iter()
.find(|position| position.instrument == "#51-OUTCOME")
.expect("no position summary");
assert_eq!(no_position.final_position_qty, "0");
assert_eq!(no_position.realized_pnl, "10");
assert_eq!(no_position.unrealized_pnl, None);
assert_eq!(no_position.net_pnl, "10");
}
#[tokio::test]
async fn action_limit_defers_batch_and_later_single_without_rejection() {
let exchange_instance =
ExchangeInstance::new(ExchangeId::new("hyperliquid"), Environment::Testnet);
let instrument = InstrumentId::new("BTC-PERP");
let mock = Arc::new(MockExchange::new());
mock.queue_place_order_error(action_limit_error(2)).await;
mock.queue_place_order_success().await;
mock.queue_place_order_error(action_limit_error(1)).await;
mock.queue_place_order_success().await;
let accepted = Arc::new(AtomicUsize::new(0));
let rejected = Arc::new(AtomicUsize::new(0));
let strategy = Box::new(DeferredPlacementProbeStrategy::new(
exchange_instance,
instrument.clone(),
accepted.clone(),
rejected.clone(),
));
let mut engine = Engine::new(EngineConfig::default());
engine.register_strategy(strategy);
engine.register_instrument(instrument_meta(&instrument));
engine.register_exchange(mock.clone() as Arc<dyn Exchange>);
let mut runner = EngineRunner::new(
engine,
RunnerConfig {
min_poll_delay_ms: 1,
quote_poll_interval_ms: 1,
cleanup_delay_ms: 0,
..Default::default()
},
);
runner.add_exchange(mock.clone() as Arc<dyn Exchange>);
runner.add_instrument(instrument);
timeout(Duration::from_secs(2), runner.run())
.await
.expect("runner should complete after deferred retries");
assert_eq!(accepted.load(Ordering::SeqCst), 3);
assert_eq!(rejected.load(Ordering::SeqCst), 0);
let placed = mock.placed_orders().await;
assert_eq!(placed.len(), 3);
assert_eq!(placed[0].client_id, ClientOrderId::new("batch-1"));
assert_eq!(placed[1].client_id, ClientOrderId::new("batch-2"));
assert_eq!(placed[2].client_id, ClientOrderId::new("later-single"));
}
}