pub mod config;
pub mod ids_generator;
pub mod inflight;
mod settlement;
use std::{
cell::RefCell,
cmp::min,
fmt::Debug,
mem,
ops::{Add, Sub},
rc::Rc,
};
use indexmap::{IndexMap, IndexSet};
use jiff::SignedDuration;
use nautilus_common::{
cache::Cache,
clock::Clock,
messages::execution::{
BatchCancelOrders, BatchModifyOrders, CancelAllOrders, CancelOrder, ModifyOrder,
},
msgbus::{self, MessagingSwitchboard},
};
use nautilus_core::{UUID4, UnixNanos, correctness::CorrectnessResult};
use nautilus_model::{
data::{
Bar, BarType, InstrumentClose, OrderBookDelta, OrderBookDeltas, OrderBookDepth10,
QuoteTick, TradeTick,
order::{BookOrder, OrderId},
},
enums::{
AccountType, AggregationSource, AggressorSide, BookAction, BookType, ContingencyType,
InstrumentCloseType, LiquiditySide, MarketStatus, MarketStatusAction, OmsType, OrderSide,
OrderStatus, OrderType, PositionSide, PriceType, RecordFlag, TimeInForce, TriggerType,
},
events::{
OrderAccepted, OrderCancelRejected, OrderCanceled, OrderEventAny, OrderExpired,
OrderFilled, OrderModifyRejected, OrderRejected, OrderSubmitted, OrderTriggered,
OrderUpdated,
},
identifiers::{
AccountId, ClientOrderId, InstrumentId, PositionId, StrategyId, TradeId, TraderId, Venue,
VenueOrderId,
},
instruments::{Instrument, InstrumentAny},
orderbook::{BookLevel, OrderBook},
orders::{MarketOrder, Order, OrderAny, OrderCore},
position::{Position, PositionReplayEvent},
types::{
Currency, Money, Price, Quantity, fixed::FIXED_PRECISION, price::PriceRaw,
quantity::QuantityRaw,
},
};
use rust_decimal::Decimal;
use ustr::Ustr;
use self::{
config::OrderMatchingEngineConfig, ids_generator::IdsGenerator, inflight::InflightOrders,
};
use crate::{
matching_core::{MatchAction, OrderMatchingCore, RestingOrder},
models::{
fee::{FeeModel, FeeModelHandle},
fill::{FillModel, FillModelHandle},
},
protection::protection_price_calculate,
trailing::trailing_stop_calculate,
};
pub struct OrderMatchingEngine {
pub venue: Venue,
pub instrument: InstrumentAny,
pub raw_id: u32,
pub book_type: BookType,
pub oms_type: OmsType,
pub account_type: AccountType,
pub market_status: MarketStatus,
pub config: OrderMatchingEngineConfig,
core: OrderMatchingCore,
clock: Rc<RefCell<dyn Clock>>,
cache: Rc<RefCell<Cache>>,
book: OrderBook,
fill_model: FillModelHandle,
fee_model: FeeModelHandle,
event_handler: Option<Rc<dyn Fn(OrderEventAny)>>,
inflight_orders: InflightOrders,
target_bid: Option<Price>,
target_ask: Option<Price>,
target_last: Option<Price>,
last_bar_bid: Option<Bar>,
last_bar_ask: Option<Bar>,
fill_at_market: bool,
execution_bar_types: IndexMap<InstrumentId, BarType>,
execution_bar_deltas: IndexMap<BarType, SignedDuration>,
account_ids: IndexMap<TraderId, AccountId>,
cached_filled_qty: IndexMap<ClientOrderId, Quantity>,
pending_order_updates: RefCell<IndexMap<ClientOrderId, Vec<OrderUpdated>>>,
pending_fills: IndexMap<TradeId, PendingFill>,
post_match_order_ids: IndexSet<ClientOrderId>,
ids_generator: IdsGenerator,
last_trade_size: Option<Quantity>,
trade_consumption: QuantityRaw,
bid_consumption: IndexMap<PriceRaw, (QuantityRaw, QuantityRaw)>,
ask_consumption: IndexMap<PriceRaw, (QuantityRaw, QuantityRaw)>,
queue_pending: IndexMap<ClientOrderId, PriceRaw>,
queue_ahead_orders: IndexMap<ClientOrderId, IndexMap<OrderId, QuantityRaw>>,
queue_ahead_total: IndexMap<ClientOrderId, (PriceRaw, QuantityRaw)>,
queue_snapshot_in_progress: bool,
queue_ids_by_price: IndexMap<PriceRaw, IndexSet<ClientOrderId>>,
queue_excess: IndexMap<ClientOrderId, QuantityRaw>,
queue_id_scratch: Vec<ClientOrderId>,
queue_pending_scratch: Vec<(ClientOrderId, PriceRaw)>,
queue_stale_scratch: Vec<ClientOrderId>,
queue_entry_scratch: Vec<(ClientOrderId, QuantityRaw, QuantityRaw)>,
prev_bid_price_raw: PriceRaw,
prev_ask_price_raw: PriceRaw,
tob_initialized: bool,
last_quote_bid: Option<Price>,
last_quote_ask: Option<Price>,
precision_mismatch_streak: u32,
instrument_close: Option<InstrumentClose>,
pending_resolution: bool,
expiration_processed: bool,
option_settlement_failed: bool,
option_settlement_warning: Option<&'static str>,
option_expiration_orders_canceled: bool,
}
impl Debug for OrderMatchingEngine {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct(stringify!(OrderMatchingEngine))
.field("venue", &self.venue)
.field("instrument", &self.instrument.id())
.finish()
}
}
impl OrderMatchingEngine {
#[expect(clippy::too_many_arguments)]
pub fn new(
instrument: InstrumentAny,
raw_id: u32,
fill_model: FillModelHandle,
fee_model: FeeModelHandle,
book_type: BookType,
oms_type: OmsType,
account_type: AccountType,
clock: Rc<RefCell<dyn Clock>>,
cache: Rc<RefCell<Cache>>,
config: OrderMatchingEngineConfig,
) -> Self {
let book = OrderBook::new(instrument.id(), book_type);
let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
core.set_fill_limit_inside_spread(Self::fill_limit_inside_spread_or_false(&fill_model));
let ids_generator = IdsGenerator::new(
instrument.id().venue,
oms_type,
raw_id,
config.use_random_ids,
config.use_position_ids,
cache.clone(),
);
Self {
venue: instrument.id().venue,
instrument,
raw_id,
fill_model,
fee_model,
event_handler: None,
inflight_orders: InflightOrders::default(),
book_type,
oms_type,
account_type,
clock,
cache,
book,
market_status: MarketStatus::Open,
config,
core,
target_bid: None,
target_ask: None,
target_last: None,
last_bar_bid: None,
last_bar_ask: None,
fill_at_market: true,
execution_bar_types: IndexMap::new(),
execution_bar_deltas: IndexMap::new(),
account_ids: IndexMap::new(),
cached_filled_qty: IndexMap::new(),
pending_order_updates: RefCell::new(IndexMap::new()),
pending_fills: IndexMap::new(),
post_match_order_ids: IndexSet::new(),
ids_generator,
last_trade_size: None,
trade_consumption: 0,
bid_consumption: IndexMap::new(),
ask_consumption: IndexMap::new(),
queue_pending: IndexMap::new(),
queue_ahead_orders: IndexMap::new(),
queue_ahead_total: IndexMap::new(),
queue_snapshot_in_progress: false,
queue_ids_by_price: IndexMap::new(),
queue_excess: IndexMap::new(),
queue_id_scratch: Vec::new(),
queue_pending_scratch: Vec::new(),
queue_stale_scratch: Vec::new(),
queue_entry_scratch: Vec::new(),
prev_bid_price_raw: 0,
prev_ask_price_raw: 0,
tob_initialized: false,
last_quote_bid: None,
last_quote_ask: None,
precision_mismatch_streak: 0,
instrument_close: None,
pending_resolution: false,
expiration_processed: false,
option_settlement_failed: false,
option_settlement_warning: None,
option_expiration_orders_canceled: false,
}
}
pub fn set_event_handler(&mut self, handler: Rc<dyn Fn(OrderEventAny)>) {
self.event_handler = Some(handler);
}
pub fn set_inflight_orders(&mut self, orders: InflightOrders) {
self.inflight_orders = orders;
}
fn dispatch_order_event(&self, event: OrderEventAny) {
if let Some(handler) = &self.event_handler {
handler(event);
} else {
let endpoint = MessagingSwitchboard::exec_engine_process();
msgbus::send_order_event(endpoint, event);
}
}
pub fn reset(&mut self) {
self.book.reset();
self.execution_bar_types.clear();
self.execution_bar_deltas.clear();
self.account_ids.clear();
self.cached_filled_qty.clear();
self.pending_order_updates.get_mut().clear();
self.pending_fills.clear();
self.post_match_order_ids.clear();
self.core.reset();
self.target_bid = None;
self.target_ask = None;
self.target_last = None;
self.last_trade_size = None;
self.trade_consumption = 0;
self.bid_consumption.clear();
self.ask_consumption.clear();
self.queue_pending.clear();
self.queue_ahead_orders.clear();
self.queue_ahead_total.clear();
self.queue_snapshot_in_progress = false;
self.queue_ids_by_price.clear();
self.queue_excess.clear();
self.queue_id_scratch.clear();
self.queue_pending_scratch.clear();
self.queue_stale_scratch.clear();
self.queue_entry_scratch.clear();
self.prev_bid_price_raw = 0;
self.prev_ask_price_raw = 0;
self.tob_initialized = false;
self.last_quote_bid = None;
self.last_quote_ask = None;
self.last_bar_bid = None;
self.last_bar_ask = None;
self.precision_mismatch_streak = 0;
self.instrument_close = None;
self.market_status = MarketStatus::Open;
self.pending_resolution = false;
self.expiration_processed = false;
self.option_settlement_failed = false;
self.option_settlement_warning = None;
self.option_expiration_orders_canceled = false;
self.fill_at_market = true;
self.ids_generator.reset();
log::info!("Reset {}", self.instrument.id());
}
fn apply_liquidity_consumption(
&mut self,
mut fills: Vec<(Price, Quantity)>,
order_side: OrderSide,
leaves_qty: Quantity,
book_prices: Option<&[Price]>,
) -> Vec<(Price, Quantity)> {
if !self.config.liquidity_consumption {
return fills;
}
let consumption = match order_side {
OrderSide::Buy => &mut self.ask_consumption,
OrderSide::Sell => &mut self.bid_consumption,
};
let mut adjusted_len = 0;
let mut remaining_qty = leaves_qty.raw();
for fill_idx in 0..fills.len() {
if remaining_qty == 0 {
break;
}
let (price, qty) = fills[fill_idx];
let book_price = book_prices
.and_then(|bp| bp.get(fill_idx).copied())
.unwrap_or(price);
let book_price_raw = book_price.raw();
let level_size = self
.book
.get_quantity_at_level(book_price, order_side, qty.precision);
let (original_size, consumed) = consumption
.entry(book_price_raw)
.or_insert((level_size.raw(), 0));
if *original_size != level_size.raw() {
*original_size = level_size.raw();
*consumed = 0;
}
let available = original_size.saturating_sub(*consumed);
if available == 0 {
continue;
}
let adjusted_qty_raw = min(min(qty.raw(), available), remaining_qty);
if adjusted_qty_raw == 0 {
continue;
}
*consumed += adjusted_qty_raw;
remaining_qty -= adjusted_qty_raw;
let adjusted_qty = Quantity::from_raw(adjusted_qty_raw, qty.precision);
fills[adjusted_len] = (price, adjusted_qty);
adjusted_len += 1;
}
fills.truncate(adjusted_len);
fills
}
fn seed_trade_consumption(
&mut self,
trade_price_raw: PriceRaw,
trade_size_raw: QuantityRaw,
trade_ts_event: UnixNanos,
aggressor_side: AggressorSide,
) {
if trade_size_raw == 0 {
return;
}
if self.book.ts_last > trade_ts_event {
return;
}
let book = &self.book;
let consumption = match aggressor_side {
AggressorSide::Buy => &mut self.ask_consumption,
AggressorSide::Sell => &mut self.bid_consumption,
AggressorSide::NoAggressor => return,
};
let mut remaining = trade_size_raw;
match aggressor_side {
AggressorSide::Buy => {
for level in book
.asks(None)
.take_while(|level| level.price.value.raw() <= trade_price_raw)
{
Self::consume_trade_level(consumption, &mut remaining, level);
if remaining == 0 {
break;
}
}
}
AggressorSide::Sell => {
for level in book
.bids(None)
.take_while(|level| level.price.value.raw() >= trade_price_raw)
{
Self::consume_trade_level(consumption, &mut remaining, level);
if remaining == 0 {
break;
}
}
}
AggressorSide::NoAggressor => unreachable!(),
}
}
fn consume_trade_level(
consumption: &mut IndexMap<PriceRaw, (QuantityRaw, QuantityRaw)>,
remaining: &mut QuantityRaw,
level: &BookLevel,
) {
let level_size = level.size_raw();
let entry = consumption
.entry(level.price.value.raw())
.or_insert((level_size, 0));
if entry.0 != level_size {
entry.0 = level_size;
entry.1 = 0;
}
let available = level_size.saturating_sub(entry.1);
let consume = min(*remaining, available);
entry.1 += consume;
*remaining -= consume;
}
pub fn set_fill_model(&mut self, fill_model: FillModelHandle) {
self.core
.set_fill_limit_inside_spread(Self::fill_limit_inside_spread_or_false(&fill_model));
self.fill_model = fill_model;
}
fn fill_limit_inside_spread_or_false(fill_model: &FillModelHandle) -> bool {
fill_model.fill_limit_inside_spread().unwrap_or_else(|e| {
log::error!("Failed to query fill model spread behavior: {e}");
false
})
}
fn snapshot_queue_position(&mut self, order: &OrderAny, price: Price) {
if !self.config.queue_position {
return;
}
let size_prec = self.instrument.size_precision();
let qty_ahead = self.book.get_quantity_at_level(
price,
OrderCore::opposite_side(order.order_side()),
size_prec,
);
let client_order_id = order.client_order_id();
self.remove_queue_position(client_order_id);
self.queue_ids_by_price
.entry(price.raw())
.or_default()
.insert(client_order_id);
if self.book_type == BookType::L1_MBP && qty_ahead.is_zero() {
let behind_bbo = match order.order_side() {
OrderSide::Buy => self.book.best_bid_price().is_some_and(|bid| price < bid),
OrderSide::Sell => self.book.best_ask_price().is_some_and(|ask| price > ask),
};
if behind_bbo {
self.queue_pending.insert(client_order_id, price.raw());
return;
}
}
self.queue_ahead_total
.insert(client_order_id, (price.raw(), qty_ahead.raw()));
if self.book_type == BookType::L3_MBO {
let orders_ahead: IndexMap<OrderId, QuantityRaw> = self
.book
.get_orders_at_level(price, OrderCore::opposite_side(order.order_side()))
.iter()
.map(|book_order| (book_order.order_id, book_order.size.raw()))
.collect();
self.queue_ahead_orders
.insert(client_order_id, orders_ahead);
}
}
fn remove_queue_position(&mut self, client_order_id: ClientOrderId) {
let pending_price = self.queue_pending.shift_remove(&client_order_id);
let ahead_price = self
.queue_ahead_total
.shift_remove(&client_order_id)
.map(|(price_raw, _)| price_raw);
self.queue_ahead_orders.shift_remove(&client_order_id);
self.queue_excess.shift_remove(&client_order_id);
for price_raw in [pending_price, ahead_price].into_iter().flatten() {
let remove_price = self
.queue_ids_by_price
.get_mut(&price_raw)
.is_some_and(|ids| {
ids.shift_remove(&client_order_id);
ids.is_empty()
});
if remove_price {
self.queue_ids_by_price.shift_remove(&price_raw);
}
}
}
fn take_queue_ids_at_price(&mut self, price_raw: PriceRaw) -> Vec<ClientOrderId> {
let mut ids = Self::take_cleared(&mut self.queue_id_scratch);
if let Some(tracked_ids) = self.queue_ids_by_price.get(&price_raw) {
ids.extend(tracked_ids.iter().copied());
}
ids
}
fn decrement_queue_on_trade(
&mut self,
price_raw: PriceRaw,
trade_size_raw: QuantityRaw,
aggressor_side: AggressorSide,
) {
if !self.config.queue_position {
return;
}
self.queue_excess.clear();
let keys = self.take_queue_ids_at_price(price_raw);
let mut entries = Self::take_cleared(&mut self.queue_entry_scratch);
let mut stale = Self::take_cleared(&mut self.queue_stale_scratch);
for client_order_id in keys.iter().copied() {
let (order_price_raw, ahead_raw) =
match self.queue_ahead_total.get(&client_order_id).copied() {
Some(v) => v,
None => continue,
};
let cache = self.cache.borrow();
let order_info = cache.order(&client_order_id).and_then(|order| {
if order.is_closed() {
return None;
}
let has_pending_updates = self
.pending_order_updates
.borrow()
.contains_key(&client_order_id);
let has_pending_fills = self
.cached_filled_qty
.get(&client_order_id)
.is_some_and(|filled_qty| *filled_qty != order.filled_qty());
let snapshot;
let order = if has_pending_updates || has_pending_fills {
snapshot = self.order_snapshot(client_order_id)?;
&snapshot
} else {
&order
};
Some((order.order_side(), order.leaves_qty().raw()))
});
drop(cache);
let Some((order_side, leaves_raw)) = order_info else {
stale.push(client_order_id);
continue;
};
if order_price_raw != price_raw || ahead_raw == 0 {
continue;
}
let should_decrement = matches!(aggressor_side, AggressorSide::NoAggressor)
|| (aggressor_side == AggressorSide::Buy && order_side == OrderSide::Sell)
|| (aggressor_side == AggressorSide::Sell && order_side == OrderSide::Buy);
if should_decrement {
entries.push((client_order_id, ahead_raw, leaves_raw));
}
}
for id in stale.drain(..) {
self.remove_queue_position(id);
}
entries.sort_by_key(|&(_, ahead, _)| ahead);
let mut remaining = trade_size_raw;
let mut prev_position: QuantityRaw = 0;
for (client_order_id, ahead_raw, leaves_raw) in &entries {
if remaining == 0 {
let new_ahead = ahead_raw.saturating_sub(trade_size_raw);
self.reduce_queue_ahead(*client_order_id, price_raw, *ahead_raw, new_ahead);
if new_ahead == 0 {
self.queue_excess.insert(*client_order_id, 0);
}
continue;
}
let gap = ahead_raw.saturating_sub(prev_position);
let queue_consumed = remaining.min(gap);
remaining -= queue_consumed;
if remaining == 0 && queue_consumed < gap {
let new_ahead = ahead_raw.saturating_sub(trade_size_raw);
self.reduce_queue_ahead(*client_order_id, price_raw, *ahead_raw, new_ahead);
continue;
}
self.reduce_queue_ahead(*client_order_id, price_raw, *ahead_raw, 0);
let excess = remaining.min(*leaves_raw);
self.queue_excess.insert(*client_order_id, excess);
remaining -= excess;
prev_position = ahead_raw + excess;
}
self.queue_id_scratch = keys;
self.queue_entry_scratch = entries;
self.queue_stale_scratch = stale;
}
fn reduce_queue_ahead(
&mut self,
client_order_id: ClientOrderId,
price_raw: PriceRaw,
ahead_raw: QuantityRaw,
new_ahead_raw: QuantityRaw,
) {
self.queue_ahead_total
.insert(client_order_id, (price_raw, new_ahead_raw));
self.consume_queue_ahead_orders(client_order_id, ahead_raw.saturating_sub(new_ahead_raw));
}
fn consume_queue_ahead_orders(
&mut self,
client_order_id: ClientOrderId,
mut amount_raw: QuantityRaw,
) {
let Some(orders_ahead) = self.queue_ahead_orders.get_mut(&client_order_id) else {
return;
};
while amount_raw > 0 {
let Some((&book_order_id, &size_raw)) = orders_ahead.get_index(0) else {
break;
};
if size_raw <= amount_raw {
orders_ahead.shift_remove(&book_order_id);
amount_raw -= size_raw;
} else {
orders_ahead.insert(book_order_id, size_raw - amount_raw);
amount_raw = 0;
}
}
}
fn determine_trade_fill_qty(&self, order: &OrderAny) -> Option<QuantityRaw> {
if !self.config.queue_position {
return Some(order.leaves_qty().raw());
}
let client_order_id = order.client_order_id();
if self.queue_pending.contains_key(&client_order_id) {
return None;
}
if let Some(&(tracked_price_raw, ahead_raw)) = self.queue_ahead_total.get(&client_order_id)
&& let Some(order_price) = order.price()
&& order_price.raw() == tracked_price_raw
&& ahead_raw > 0
{
return None;
}
let leaves_raw = order.leaves_qty().raw();
if leaves_raw == 0 {
return None;
}
let mut available_raw = leaves_raw;
if let Some(trade_size) = self.last_trade_size {
let remaining = trade_size.raw().saturating_sub(self.trade_consumption);
available_raw = available_raw.min(remaining);
if let Some(&excess_raw) = self.queue_excess.get(&client_order_id) {
if excess_raw == 0 {
return None;
}
available_raw = available_raw.min(excess_raw);
}
}
if available_raw == 0 {
return None;
}
Some(available_raw)
}
fn rebase_queue_positions(&mut self) {
if !self.config.queue_position {
return;
}
let tracked: Vec<_> = self
.queue_ahead_total
.iter()
.map(|(&client_order_id, &(price_raw, ahead_raw))| {
(client_order_id, price_raw, ahead_raw)
})
.collect();
let mut stale = Self::take_cleared(&mut self.queue_stale_scratch);
let size_precision = self.instrument.size_precision();
let price_precision = self.instrument.price_precision();
for (client_order_id, price_raw, ahead_raw) in tracked {
let order_side = self
.cache
.borrow()
.order(&client_order_id)
.and_then(|order| {
if order.is_closed() {
None
} else {
Some(order.order_side())
}
});
let Some(order_side) = order_side else {
stale.push(client_order_id);
continue;
};
let price = Price::from_raw(price_raw, price_precision);
let visible_raw = self
.book
.get_quantity_at_level(price, OrderCore::opposite_side(order_side), size_precision)
.raw();
let rebased_raw = ahead_raw.min(visible_raw);
if self.book_type == BookType::L3_MBO {
let previous_orders = self
.queue_ahead_orders
.get(&client_order_id)
.cloned()
.unwrap_or_default();
let mut orders_ahead = IndexMap::new();
let mut total_raw = 0;
for book_order in self
.book
.get_orders_at_level(price, OrderCore::opposite_side(order_side))
{
if !previous_orders.contains_key(&book_order.order_id) {
continue;
}
let previous_size_raw = previous_orders[&book_order.order_id];
let size_raw = previous_size_raw.min(book_order.size.raw());
orders_ahead.insert(book_order.order_id, size_raw);
total_raw += size_raw;
}
self.queue_ahead_orders
.insert(client_order_id, orders_ahead);
self.queue_ahead_total
.insert(client_order_id, (price_raw, total_raw));
} else {
self.queue_ahead_total
.insert(client_order_id, (price_raw, rebased_raw));
}
}
for client_order_id in stale.drain(..) {
self.remove_queue_position(client_order_id);
}
self.queue_stale_scratch = stale;
}
fn adjust_queue_for_delta(&mut self, delta: &OrderBookDelta) {
if delta.action == BookAction::Delete {
if self.is_order_granular_delta(delta.flags) {
self.advance_l3_queue_on_delete(delta.order.order_id);
} else {
self.clear_queue_on_delete(delta.order.price.raw(), delta.order.side);
}
} else if delta.action == BookAction::Update {
if self.is_order_granular_delta(delta.flags) {
self.adjust_l3_queue_on_update(&delta.order);
} else {
self.cap_queue_ahead(
delta.order.price.raw(),
delta.order.size.raw(),
delta.order.side,
);
}
}
}
fn clear_queue_on_delete(
&mut self,
deleted_price_raw: PriceRaw,
deleted_side: Option<OrderSide>,
) {
let keys = self.take_queue_ids_at_price(deleted_price_raw);
for client_order_id in keys.iter().copied() {
if let Some(&(order_price_raw, ahead_raw)) =
self.queue_ahead_total.get(&client_order_id)
&& order_price_raw == deleted_price_raw
{
let matches_side = self
.cache
.borrow()
.order(&client_order_id)
.is_some_and(|o| Some(o.order_side()) == deleted_side);
if matches_side {
self.reduce_queue_ahead(client_order_id, order_price_raw, ahead_raw, 0);
}
}
}
self.queue_id_scratch = keys;
}
fn is_order_granular_delta(&self, flags: u8) -> bool {
self.book_type == BookType::L3_MBO
&& !RecordFlag::F_TOB.matches(flags)
&& !RecordFlag::F_MBP.matches(flags)
}
fn advance_l3_queue_on_delete(&mut self, book_order_id: OrderId) {
for (client_order_id, orders_ahead) in &mut self.queue_ahead_orders {
let Some(size_raw) = orders_ahead.shift_remove(&book_order_id) else {
continue;
};
if let Some((_, ahead_raw)) = self.queue_ahead_total.get_mut(client_order_id) {
*ahead_raw = ahead_raw.saturating_sub(size_raw);
}
}
}
fn adjust_l3_queue_on_update(&mut self, book_order: &BookOrder) {
for (client_order_id, orders_ahead) in &mut self.queue_ahead_orders {
let Some(&tracked_size_raw) = orders_ahead.get(&book_order.order_id) else {
continue;
};
let Some((tracked_price_raw, ahead_raw)) =
self.queue_ahead_total.get_mut(client_order_id)
else {
continue;
};
if book_order.price.raw() != *tracked_price_raw {
*ahead_raw = ahead_raw.saturating_sub(tracked_size_raw);
orders_ahead.shift_remove(&book_order.order_id);
} else if book_order.size.raw() < tracked_size_raw {
*ahead_raw = ahead_raw.saturating_sub(tracked_size_raw - book_order.size.raw());
orders_ahead.insert(book_order.order_id, book_order.size.raw());
} else if book_order.size.raw() > tracked_size_raw {
*ahead_raw = ahead_raw.saturating_add(book_order.size.raw() - tracked_size_raw);
orders_ahead.insert(book_order.order_id, book_order.size.raw());
}
}
}
fn cap_queue_ahead(
&mut self,
price_raw: PriceRaw,
size_raw: QuantityRaw,
order_side: Option<OrderSide>,
) {
let keys = self.take_queue_ids_at_price(price_raw);
let mut stale = Self::take_cleared(&mut self.queue_stale_scratch);
for client_order_id in keys.iter().copied() {
let (order_price_raw, ahead_raw) =
match self.queue_ahead_total.get(&client_order_id).copied() {
Some(v) => v,
None => continue,
};
if order_price_raw != price_raw || ahead_raw <= size_raw {
continue;
}
let cache = self.cache.borrow();
let order_info = cache.order(&client_order_id).and_then(|order| {
if order.is_closed() {
None
} else {
Some(order.order_side())
}
});
drop(cache);
let Some(side) = order_info else {
stale.push(client_order_id);
continue;
};
if Some(side) != order_side {
continue;
}
self.reduce_queue_ahead(client_order_id, order_price_raw, ahead_raw, size_raw);
}
for id in stale.drain(..) {
self.remove_queue_position(id);
}
self.queue_id_scratch = keys;
self.queue_stale_scratch = stale;
}
fn seed_tob_baseline(&mut self) {
let bid = self.book.best_bid_price();
let ask = self.book.best_ask_price();
self.prev_bid_price_raw = bid.map_or(0, |p| p.raw());
self.prev_ask_price_raw = ask.map_or(0, |p| p.raw());
self.tob_initialized = bid.is_some() || ask.is_some();
}
fn decrement_l1_queue_on_quote(
&mut self,
bid_price_raw: PriceRaw,
bid_size_raw: QuantityRaw,
ask_price_raw: PriceRaw,
ask_size_raw: QuantityRaw,
) {
if !self.config.queue_position {
return;
}
if self.tob_initialized {
if bid_price_raw < self.prev_bid_price_raw {
self.adjust_l1_queue_on_price_move(bid_price_raw, bid_size_raw, OrderSide::Buy);
}
if ask_price_raw > self.prev_ask_price_raw {
self.adjust_l1_queue_on_price_move(ask_price_raw, ask_size_raw, OrderSide::Sell);
}
}
self.resolve_pending_l1_snapshots(bid_price_raw, bid_size_raw, ask_price_raw, ask_size_raw);
}
fn adjust_l1_queue_on_price_move(
&mut self,
new_price_raw: PriceRaw,
new_size_raw: QuantityRaw,
order_side: OrderSide,
) {
let mut keys = Self::take_cleared(&mut self.queue_id_scratch);
keys.extend(self.queue_ahead_total.keys().copied());
let mut stale = Self::take_cleared(&mut self.queue_stale_scratch);
for client_order_id in keys.iter().copied() {
let Some(&(order_price_raw, ahead_raw)) = self.queue_ahead_total.get(&client_order_id)
else {
continue;
};
let cache = self.cache.borrow();
let order_info = cache.order(&client_order_id).and_then(|order| {
if order.is_closed() {
None
} else {
Some(order.order_side())
}
});
drop(cache);
let Some(side) = order_info else {
stale.push(client_order_id);
continue;
};
if side != order_side {
continue;
}
let crossed = match order_side {
OrderSide::Buy => order_price_raw > new_price_raw,
_ => order_price_raw < new_price_raw,
};
if crossed {
self.queue_ahead_total
.insert(client_order_id, (order_price_raw, 0));
} else if order_price_raw == new_price_raw && ahead_raw > new_size_raw {
self.queue_ahead_total
.insert(client_order_id, (order_price_raw, new_size_raw));
}
}
for id in stale.drain(..) {
self.remove_queue_position(id);
}
let mut pending = Self::take_cleared(&mut self.queue_pending_scratch);
pending.extend(
self.queue_pending
.iter()
.map(|(&client_order_id, &price_raw)| (client_order_id, price_raw)),
);
for (client_order_id, order_price_raw) in pending.iter().copied() {
let cache = self.cache.borrow();
let order_info = cache.order(&client_order_id).and_then(|order| {
if order.is_closed() {
None
} else {
Some(order.order_side())
}
});
drop(cache);
let Some(side) = order_info else {
stale.push(client_order_id);
continue;
};
if side != order_side {
continue;
}
let crossed = match order_side {
OrderSide::Buy => order_price_raw > new_price_raw,
_ => order_price_raw < new_price_raw,
};
if crossed {
self.queue_pending.shift_remove(&client_order_id);
self.queue_ahead_total
.insert(client_order_id, (order_price_raw, 0));
} else if order_price_raw == new_price_raw {
self.queue_pending.shift_remove(&client_order_id);
self.queue_ahead_total
.insert(client_order_id, (order_price_raw, new_size_raw));
}
}
for id in stale.drain(..) {
self.remove_queue_position(id);
}
self.queue_id_scratch = keys;
self.queue_pending_scratch = pending;
self.queue_stale_scratch = stale;
}
fn resolve_pending_l1_snapshots(
&mut self,
bid_price_raw: PriceRaw,
bid_size_raw: QuantityRaw,
ask_price_raw: PriceRaw,
ask_size_raw: QuantityRaw,
) {
let mut keys = self.take_queue_ids_at_price(bid_price_raw);
if ask_price_raw != bid_price_raw
&& let Some(ask_ids) = self.queue_ids_by_price.get(&ask_price_raw)
{
keys.extend(ask_ids.iter().copied());
}
let mut stale = Self::take_cleared(&mut self.queue_stale_scratch);
for client_order_id in keys.iter().copied() {
let Some(&order_price_raw) = self.queue_pending.get(&client_order_id) else {
continue;
};
let cache = self.cache.borrow();
let order_info = cache.order(&client_order_id).and_then(|order| {
if order.is_closed() {
None
} else {
Some(order.order_side())
}
});
drop(cache);
let Some(side) = order_info else {
stale.push(client_order_id);
continue;
};
let matched_size = match side {
OrderSide::Buy if order_price_raw == bid_price_raw => Some(bid_size_raw),
OrderSide::Sell if order_price_raw == ask_price_raw => Some(ask_size_raw),
_ => None,
};
if let Some(size) = matched_size {
self.queue_pending.shift_remove(&client_order_id);
self.queue_ahead_total
.insert(client_order_id, (order_price_raw, size));
}
}
for id in stale.drain(..) {
self.remove_queue_position(id);
}
self.queue_id_scratch = keys;
self.queue_stale_scratch = stale;
}
fn resolve_pending_on_trade(&mut self, trade_price_raw: PriceRaw) {
let mut keys = Self::take_cleared(&mut self.queue_id_scratch);
keys.extend(self.queue_pending.keys().copied());
let mut stale = Self::take_cleared(&mut self.queue_stale_scratch);
for client_order_id in keys.iter().copied() {
let Some(&order_price_raw) = self.queue_pending.get(&client_order_id) else {
continue;
};
let cache = self.cache.borrow();
let order_side = cache.order(&client_order_id).and_then(|order| {
if order.is_closed() {
None
} else {
Some(order.order_side())
}
});
drop(cache);
let Some(side) = order_side else {
stale.push(client_order_id);
continue;
};
let crossed = match side {
OrderSide::Buy => trade_price_raw < order_price_raw,
OrderSide::Sell => trade_price_raw > order_price_raw,
};
if crossed {
self.queue_pending.shift_remove(&client_order_id);
self.queue_ahead_total
.insert(client_order_id, (order_price_raw, 0));
}
}
for id in stale.drain(..) {
self.remove_queue_position(id);
}
self.queue_id_scratch = keys;
self.queue_stale_scratch = stale;
}
fn take_cleared<T>(buf: &mut Vec<T>) -> Vec<T> {
let mut items = mem::take(buf);
items.clear();
items
}
#[must_use]
pub fn best_bid_price(&self) -> Option<Price> {
self.book.best_bid_price()
}
#[must_use]
pub fn best_ask_price(&self) -> Option<Price> {
self.book.best_ask_price()
}
#[must_use]
pub const fn get_book(&self) -> &OrderBook {
&self.book
}
#[must_use]
pub fn get_open_bid_orders(&self) -> Vec<RestingOrder> {
self.core.get_orders_bid()
}
#[must_use]
pub fn get_open_ask_orders(&self) -> Vec<RestingOrder> {
self.core.get_orders_ask()
}
#[must_use]
pub fn get_open_orders(&self) -> Vec<RestingOrder> {
self.core.get_orders()
}
#[must_use]
pub fn order_exists(&self, client_order_id: ClientOrderId) -> bool {
self.core.order_exists(client_order_id)
}
#[must_use]
pub fn cached_filled_qty_len(&self) -> usize {
self.cached_filled_qty.len()
}
#[must_use]
pub const fn get_core(&self) -> &OrderMatchingCore {
&self.core
}
pub fn set_fill_at_market(&mut self, value: bool) {
self.fill_at_market = value;
}
pub fn update_instrument(&mut self, instrument: InstrumentAny) -> anyhow::Result<()> {
if instrument.id() != self.instrument.id() {
anyhow::bail!(
"Cannot update instrument {} with {}",
self.instrument.id(),
instrument.id()
);
}
let changed = instrument.price_increment() != self.instrument.price_increment()
|| instrument.price_precision() != self.instrument.price_precision()
|| instrument.size_precision() != self.instrument.size_precision();
if changed {
self.core
.update_price_increment(instrument.price_increment());
self.book.reset();
self.trade_consumption = 0;
self.bid_consumption.clear();
self.ask_consumption.clear();
self.queue_pending.clear();
self.queue_ahead_orders.clear();
self.queue_ahead_total.clear();
self.queue_ids_by_price.clear();
self.queue_excess.clear();
self.prev_bid_price_raw = 0;
self.prev_ask_price_raw = 0;
self.tob_initialized = false;
self.last_quote_bid = None;
self.last_quote_ask = None;
self.precision_mismatch_streak = 0;
self.target_bid = None;
self.target_ask = None;
self.target_last = None;
self.last_bar_bid = None;
self.last_bar_ask = None;
self.core.bid = None;
self.core.ask = None;
self.core.last = None;
log::info!(
"Updated instrument {} (price_precision={} size_precision={})",
instrument.id(),
instrument.price_precision(),
instrument.size_precision()
);
}
self.instrument = instrument;
if changed {
self.drop_incompatible_core_orders();
}
Ok(())
}
fn check_price_precision(&self, actual: u8, field: &str) -> anyhow::Result<()> {
let expected = self.instrument.price_precision();
if actual != expected {
anyhow::bail!(
"Invalid {field} precision {actual}, expected {expected} for {}",
self.instrument.id()
);
}
Ok(())
}
fn check_size_precision(&self, actual: u8, field: &str) -> anyhow::Result<()> {
let expected = self.instrument.size_precision();
if actual != expected {
anyhow::bail!(
"Invalid {field} precision {actual}, expected {expected} for {}",
self.instrument.id()
);
}
Ok(())
}
fn log_precision_mismatch(
&mut self,
data_type: &str,
instrument_id: InstrumentId,
err: &anyhow::Error,
) {
self.precision_mismatch_streak = self.precision_mismatch_streak.saturating_add(1);
let streak = self.precision_mismatch_streak;
if streak <= 3 || streak.is_multiple_of(100) {
log::warn!(
"Skipping {data_type} for {instrument_id}: {err} \
(consecutive_precision_mismatches={streak})"
);
}
if streak == 20 {
log::error!(
"Precision mismatches reached {streak} consecutive events for \
{instrument_id}; check instrument update flow and upstream market data"
);
}
}
fn drop_incompatible_core_orders(&mut self) {
let client_order_ids: Vec<ClientOrderId> = self
.core
.iter_orders()
.filter(|order| {
!self.resting_order_matches_current_instrument(order)
|| !self.cached_order_matches_current_instrument(order.client_order_id)
})
.map(|order| order.client_order_id)
.collect();
for client_order_id in client_order_ids {
let order = self
.cache
.borrow()
.order(&client_order_id)
.map(|o| o.clone());
if let Some(order) = order
&& (order.is_inflight() || order.is_open())
{
log::warn!(
"Canceling order {client_order_id} after instrument update: \
price, trigger price, or quantity is not compatible with {}",
self.instrument.id()
);
self.cancel_order(&order, None);
} else {
self.delete_core_order(client_order_id);
self.cached_filled_qty.swap_remove(&client_order_id);
}
}
}
fn cached_order_matches_current_instrument(&self, client_order_id: ClientOrderId) -> bool {
self.cache
.borrow()
.order(&client_order_id)
.is_none_or(|order| {
Self::quantity_matches_precision(order.quantity(), self.instrument.size_precision())
})
}
fn resting_order_matches_current_instrument(&self, order: &RestingOrder) -> bool {
order
.limit_price
.is_none_or(|price| self.price_matches_current_instrument(price))
&& order
.trigger_price
.is_none_or(|price| self.price_matches_current_instrument(price))
}
fn price_matches_current_instrument(&self, price: Price) -> bool {
Self::price_matches_precision(price, self.instrument.price_precision())
&& Self::price_matches_tick(price, self.instrument.price_increment())
}
fn price_matches_precision(price: Price, precision: u8) -> bool {
let precision_diff = FIXED_PRECISION.saturating_sub(precision);
let scale = PriceRaw::pow(10, u32::from(precision_diff));
price.raw() % scale == 0
}
fn price_matches_tick(price: Price, increment: Price) -> bool {
let increment_raw = increment.raw().abs();
increment_raw == 0 || price.raw() % increment_raw == 0
}
fn quantity_matches_precision(quantity: Quantity, precision: u8) -> bool {
let precision_diff = FIXED_PRECISION.saturating_sub(precision);
let scale = QuantityRaw::pow(10, u32::from(precision_diff));
quantity.raw().is_multiple_of(scale)
}
fn normalize_price_for_current_instrument(&self, price: Price) -> Option<Price> {
if !self.price_matches_current_instrument(price) {
return None;
}
Some(Price::from_raw(
price.raw(),
self.instrument.price_precision(),
))
}
fn normalize_quantity_for_current_instrument(&self, quantity: Quantity) -> Option<Quantity> {
let precision = self.instrument.size_precision();
if !Self::quantity_matches_precision(quantity, precision) {
return None;
}
Some(Quantity::from_raw(quantity.raw(), precision))
}
pub fn process_order_book_delta(&mut self, delta: &OrderBookDelta) -> anyhow::Result<()> {
log::debug!("Processing {delta}");
if matches!(delta.action, BookAction::Add | BookAction::Update) {
self.check_price_precision(delta.order.price.precision, "delta order price")?;
self.check_size_precision(delta.order.size.precision, "delta order size")?;
}
if self.book_type == BookType::L1_MBP {
self.iterate(delta.ts_init, AggressorSide::NoAggressor);
return Ok(());
}
self.book.apply_delta(delta)?;
let is_snapshot = RecordFlag::F_SNAPSHOT.matches(delta.flags);
let is_last = RecordFlag::F_LAST.matches(delta.flags);
let is_clear = delta.action == BookAction::Clear;
let snapshot_complete = is_last && (is_snapshot || self.queue_snapshot_in_progress);
if self.config.queue_position {
if is_snapshot && !is_last {
self.queue_snapshot_in_progress = true;
}
if snapshot_complete {
self.queue_snapshot_in_progress = false;
self.rebase_queue_positions();
} else if is_clear && !is_snapshot {
self.rebase_queue_positions();
} else if !self.queue_snapshot_in_progress {
self.adjust_queue_for_delta(delta);
}
}
if self.config.queue_position && (snapshot_complete || (is_clear && !is_snapshot)) {
self.seed_tob_baseline();
}
self.iterate(delta.ts_init, AggressorSide::NoAggressor);
Ok(())
}
pub fn process_order_book_deltas(&mut self, deltas: &OrderBookDeltas) -> anyhow::Result<()> {
log::debug!("Processing {deltas}");
for delta in &deltas.deltas {
if matches!(delta.action, BookAction::Add | BookAction::Update) {
self.check_price_precision(delta.order.price.precision, "delta order price")?;
self.check_size_precision(delta.order.size.precision, "delta order size")?;
}
}
if self.book_type == BookType::L1_MBP {
self.iterate(deltas.ts_init, AggressorSide::NoAggressor);
return Ok(());
}
self.book.apply_deltas(deltas)?;
let mut has_snapshot_or_clear = false;
if self.config.queue_position {
for delta in &deltas.deltas {
if RecordFlag::F_SNAPSHOT.matches(delta.flags) || delta.action == BookAction::Clear
{
has_snapshot_or_clear = true;
break;
}
self.adjust_queue_for_delta(delta);
}
}
if self.config.queue_position && has_snapshot_or_clear {
self.queue_snapshot_in_progress = false;
self.rebase_queue_positions();
self.seed_tob_baseline();
}
self.iterate(deltas.ts_init, AggressorSide::NoAggressor);
Ok(())
}
pub fn process_order_book_depth10(&mut self, depth: &OrderBookDepth10) -> anyhow::Result<()> {
log::debug!("Processing OrderBookDepth10 for {}", depth.instrument_id);
for order in &depth.bids {
if order.side.is_none() || !order.size.is_positive() {
continue;
}
self.check_price_precision(order.price.precision, "bid price")?;
self.check_size_precision(order.size.precision, "bid size")?;
}
for order in &depth.asks {
if order.side.is_none() || !order.size.is_positive() {
continue;
}
self.check_price_precision(order.price.precision, "ask price")?;
self.check_size_precision(order.size.precision, "ask size")?;
}
let top_bid = Self::first_valid_depth_order(&depth.bids, OrderSide::Buy);
let top_ask = Self::first_valid_depth_order(&depth.asks, OrderSide::Sell);
if self.book_type == BookType::L1_MBP {
let quote = QuoteTick::new(
depth.instrument_id,
Self::depth_quote_price(top_bid, self.instrument.price_precision()),
Self::depth_quote_price(top_ask, self.instrument.price_precision()),
Self::depth_quote_size(top_bid, self.instrument.size_precision()),
Self::depth_quote_size(top_ask, self.instrument.size_precision()),
depth.ts_event,
depth.ts_init,
);
self.book.update_quote_tick("e)?;
self.last_quote_bid = top_bid.map(|order| order.price);
self.last_quote_ask = top_ask.map(|order| order.price);
} else {
self.book.apply_depth(depth)?;
}
if self.config.queue_position {
self.rebase_queue_positions();
let bid_price_raw = top_bid.map_or(0, |order| order.price.raw());
let bid_size_raw = top_bid.map_or(0, |order| order.size.raw());
let ask_price_raw = top_ask.map_or(0, |order| order.price.raw());
let ask_size_raw = top_ask.map_or(0, |order| order.size.raw());
self.decrement_l1_queue_on_quote(
bid_price_raw,
bid_size_raw,
ask_price_raw,
ask_size_raw,
);
self.prev_bid_price_raw = bid_price_raw;
self.prev_ask_price_raw = ask_price_raw;
self.tob_initialized = true;
}
self.iterate(depth.ts_init, AggressorSide::NoAggressor);
Ok(())
}
fn first_valid_depth_order(orders: &[BookOrder], side: OrderSide) -> Option<BookOrder> {
orders
.iter()
.copied()
.find(|order| order.side == Some(side) && order.size.is_positive())
}
fn depth_quote_price(order: Option<BookOrder>, price_precision: u8) -> Price {
order.map_or_else(|| Price::zero(price_precision), |order| order.price)
}
fn depth_quote_size(order: Option<BookOrder>, size_precision: u8) -> Quantity {
order.map_or_else(|| Quantity::zero(size_precision), |order| order.size)
}
pub fn process_quote_tick(&mut self, quote: &QuoteTick) {
log::debug!("Processing {quote}");
if let Err(e) = self.check_price_precision(quote.bid_price.precision, "bid_price") {
self.log_precision_mismatch("quote tick", quote.instrument_id, &e);
return;
}
if let Err(e) = self.check_price_precision(quote.ask_price.precision, "ask_price") {
self.log_precision_mismatch("quote tick", quote.instrument_id, &e);
return;
}
if let Err(e) = self.check_size_precision(quote.bid_size.precision, "bid_size") {
self.log_precision_mismatch("quote tick", quote.instrument_id, &e);
return;
}
if let Err(e) = self.check_size_precision(quote.ask_size.precision, "ask_size") {
self.log_precision_mismatch("quote tick", quote.instrument_id, &e);
return;
}
self.precision_mismatch_streak = 0;
if self.book_type == BookType::L1_MBP {
if quote.ts_event < self.book.ts_last {
log::warn!(
"Skipping stale quote: ts_event {} < book.ts_last {} for {}",
quote.ts_event,
self.book.ts_last,
self.book.instrument_id,
);
self.iterate(quote.ts_init, AggressorSide::NoAggressor);
return;
}
if !self.update_quote_tick_or_skip(quote, "quote tick") {
return;
}
if self.config.queue_position {
self.decrement_l1_queue_on_quote(
quote.bid_price.raw(),
quote.bid_size.raw(),
quote.ask_price.raw(),
quote.ask_size.raw(),
);
self.prev_bid_price_raw = quote.bid_price.raw();
self.prev_ask_price_raw = quote.ask_price.raw();
self.tob_initialized = true;
}
self.last_quote_bid = Some(quote.bid_price);
self.last_quote_ask = Some(quote.ask_price);
}
self.iterate(quote.ts_init, AggressorSide::NoAggressor);
}
pub fn process_bar(&mut self, bar: &Bar) {
log::debug!("Processing {bar}");
debug_assert!(
bar.high >= bar.open
&& bar.high >= bar.low
&& bar.high >= bar.close
&& bar.low <= bar.open
&& bar.low <= bar.close,
"OHLC invariant violated for {bar}"
);
if !self.config.bar_execution || self.book_type != BookType::L1_MBP {
return;
}
let bar_type = bar.bar_type;
if bar_type.aggregation_source() == AggregationSource::Internal {
return;
}
if let Err(e) = self.check_price_precision(bar.open.precision, "bar open") {
self.log_precision_mismatch("bar", bar.instrument_id(), &e);
return;
}
if let Err(e) = self.check_price_precision(bar.high.precision, "bar high") {
self.log_precision_mismatch("bar", bar.instrument_id(), &e);
return;
}
if let Err(e) = self.check_price_precision(bar.low.precision, "bar low") {
self.log_precision_mismatch("bar", bar.instrument_id(), &e);
return;
}
if let Err(e) = self.check_price_precision(bar.close.precision, "bar close") {
self.log_precision_mismatch("bar", bar.instrument_id(), &e);
return;
}
if let Err(e) = self.check_size_precision(bar.volume.precision, "bar volume") {
self.log_precision_mismatch("bar", bar.instrument_id(), &e);
return;
}
self.precision_mismatch_streak = 0;
let price_type = bar_type.spec().price_type;
if price_type == PriceType::Mark {
log::warn!(
"Cannot process bar for {} with `PriceType::Mark`, mark price bars are not supported for bar execution",
bar.instrument_id(),
);
return;
}
let execution_bar_type =
if let Some(execution_bar_type) = self.execution_bar_types.get(&bar.instrument_id()) {
execution_bar_type.to_owned()
} else {
self.execution_bar_types
.insert(bar.instrument_id(), bar_type);
self.execution_bar_deltas
.insert(bar_type, bar_type.spec().timedelta());
bar_type
};
if execution_bar_type != bar_type {
let mut bar_type_timedelta = self.execution_bar_deltas.get(&bar_type).copied();
if bar_type_timedelta.is_none() {
bar_type_timedelta = Some(bar_type.spec().timedelta());
self.execution_bar_deltas
.insert(bar_type, bar_type_timedelta.unwrap());
}
if self.execution_bar_deltas.get(&execution_bar_type).unwrap()
>= &bar_type_timedelta.unwrap()
{
self.execution_bar_types
.insert(bar_type.instrument_id(), bar_type);
} else {
return;
}
}
match price_type {
PriceType::Last | PriceType::Mid => self.process_trade_ticks_from_bar(bar),
PriceType::Bid => {
self.last_bar_bid = Some(bar.to_owned());
self.process_quote_ticks_from_bar();
}
PriceType::Ask => {
self.last_bar_ask = Some(bar.to_owned());
self.process_quote_ticks_from_bar();
}
PriceType::Mark => {
unreachable!("PriceType::Mark bars return before execution bar state updates")
}
}
}
fn process_trade_ticks_from_bar(&mut self, bar: &Bar) {
let sizes = BarTickSizes::from_volume(bar.volume, self.instrument.size_increment());
let aggressor_side = if self.core.last.is_none_or(|last| bar.open > last) {
AggressorSide::Buy
} else {
AggressorSide::Sell
};
if self.core.last.is_none() {
self.fill_at_market = true;
if !self.process_bar_trade_tick(
bar,
bar.open,
sizes.open,
aggressor_side,
"bar open trade tick",
) {
return;
}
self.core.set_last_raw(bar.open);
} else if self.core.last.is_some_and(|last| bar.open != last) {
self.fill_at_market = true;
if !self.process_bar_trade_tick(
bar,
bar.open,
sizes.open,
aggressor_side,
"bar gap-open trade tick",
) {
return;
}
self.core.set_last_raw(bar.open);
}
let high_first = self.bar_high_first(bar);
if high_first {
self.process_bar_high(bar, sizes.high);
self.process_bar_low(bar, sizes.low);
} else {
self.process_bar_low(bar, sizes.low);
self.process_bar_high(bar, sizes.high);
}
if self.core.last.is_some_and(|last| bar.close != last) {
self.fill_at_market = false;
let aggressor_side = if bar.close > self.core.last.unwrap() {
AggressorSide::Buy
} else {
AggressorSide::Sell
};
if !self.process_bar_trade_tick(
bar,
bar.close,
sizes.close,
aggressor_side,
"bar close trade tick",
) {
return;
}
self.core.set_last_raw(bar.close);
}
self.fill_at_market = true;
}
fn process_bar_high(&mut self, bar: &Bar, size: Quantity) {
if self.core.last.is_some_and(|last| bar.high > last) {
self.fill_at_market = false;
if !self.process_bar_trade_tick(
bar,
bar.high,
size,
AggressorSide::Buy,
"bar high trade tick",
) {
return;
}
self.core.set_last_raw(bar.high);
}
}
fn process_bar_low(&mut self, bar: &Bar, size: Quantity) {
if self.core.last.is_some_and(|last| bar.low < last) {
self.fill_at_market = false;
if !self.process_bar_trade_tick(
bar,
bar.low,
size,
AggressorSide::Sell,
"bar low trade tick",
) {
return;
}
self.core.set_last_raw(bar.low);
}
}
fn process_bar_trade_tick(
&mut self,
bar: &Bar,
price: Price,
size: Quantity,
aggressor_side: AggressorSide,
context: &str,
) -> bool {
if size.is_zero() {
return true;
}
let trade_tick = TradeTick::new(
bar.instrument_id(),
price,
size,
aggressor_side,
self.ids_generator.generate_trade_id(bar.ts_init),
bar.ts_init,
bar.ts_init,
);
if !self.update_trade_tick_or_skip(&trade_tick, context) {
return false;
}
self.iterate(trade_tick.ts_init, AggressorSide::NoAggressor);
true
}
fn process_quote_ticks_from_bar(&mut self) {
if self.last_bar_bid.is_none()
|| self.last_bar_ask.is_none()
|| self.last_bar_bid.unwrap().ts_init != self.last_bar_ask.unwrap().ts_init
{
return;
}
let bid_bar = self.last_bar_bid.unwrap();
let ask_bar = self.last_bar_ask.unwrap();
let size_increment = self.instrument.size_increment();
let bid_sizes = BarTickSizes::from_volume(bid_bar.volume, size_increment);
let ask_sizes = BarTickSizes::from_volume(ask_bar.volume, size_increment);
let mut has_current_bid = false;
let mut has_current_ask = false;
let mut quote_tick = QuoteTick::new(
self.book.instrument_id,
bid_bar.open,
ask_bar.open,
bid_sizes.open,
ask_sizes.open,
bid_bar.ts_init,
bid_bar.ts_init,
);
self.fill_at_market = true;
if !self.process_bar_quote_tick(
"e_tick,
"bar open quote tick",
&mut has_current_bid,
&mut has_current_ask,
) {
return;
}
let high_first = self.bar_high_first(&bid_bar);
let high_leg = (
bid_bar.high,
ask_bar.high,
bid_sizes.high,
ask_sizes.high,
"bar high quote tick",
);
let low_leg = (
bid_bar.low,
ask_bar.low,
bid_sizes.low,
ask_sizes.low,
"bar low quote tick",
);
let legs = if high_first {
[high_leg, low_leg]
} else {
[low_leg, high_leg]
};
for (bid_price, ask_price, bid_size, ask_size, context) in legs {
self.fill_at_market = false;
quote_tick.bid_price = bid_price;
quote_tick.ask_price = ask_price;
quote_tick.bid_size = bid_size;
quote_tick.ask_size = ask_size;
if !self.process_bar_quote_tick(
"e_tick,
context,
&mut has_current_bid,
&mut has_current_ask,
) {
return;
}
}
self.fill_at_market = false;
quote_tick.bid_price = bid_bar.close;
quote_tick.ask_price = ask_bar.close;
quote_tick.bid_size = bid_sizes.close;
quote_tick.ask_size = ask_sizes.close;
if !self.process_bar_quote_tick(
"e_tick,
"bar close quote tick",
&mut has_current_bid,
&mut has_current_ask,
) {
return;
}
self.last_bar_bid = None;
self.last_bar_ask = None;
self.fill_at_market = true;
}
fn process_bar_quote_tick(
&mut self,
quote: &QuoteTick,
context: &str,
has_current_bid: &mut bool,
has_current_ask: &mut bool,
) -> bool {
let has_bid_size = quote.bid_size.non_zero();
let has_ask_size = quote.ask_size.non_zero();
let mut book_changed = false;
let mut bid_cleared = false;
let mut ask_cleared = false;
match (has_bid_size, has_ask_size) {
(true, true) => {
if !self.update_quote_tick_or_skip(quote, context) {
return false;
}
*has_current_bid = true;
*has_current_ask = true;
book_changed = true;
}
_ => {
if has_bid_size {
self.update_bar_quote_bid(quote);
*has_current_bid = true;
book_changed = true;
} else if !*has_current_bid {
self.clear_bar_quote_bid(quote);
*has_current_bid = true;
book_changed = true;
bid_cleared = true;
}
if has_ask_size {
self.update_bar_quote_ask(quote);
*has_current_ask = true;
book_changed = true;
} else if !*has_current_ask {
self.clear_bar_quote_ask(quote);
*has_current_ask = true;
book_changed = true;
ask_cleared = true;
}
}
}
if book_changed
&& let (Some(best_bid), Some(best_ask)) =
(self.book.best_bid_price(), self.book.best_ask_price())
&& best_bid > best_ask
{
if has_bid_size && !has_ask_size {
self.clear_bar_quote_ask(quote);
ask_cleared = true;
} else if has_ask_size && !has_bid_size {
self.clear_bar_quote_bid(quote);
bid_cleared = true;
}
}
if has_bid_size {
self.last_quote_bid = Some(quote.bid_price);
} else if bid_cleared {
self.last_quote_bid = None;
}
if has_ask_size {
self.last_quote_ask = Some(quote.ask_price);
} else if ask_cleared {
self.last_quote_ask = None;
}
if !book_changed {
return true;
}
self.iterate(quote.ts_init, AggressorSide::NoAggressor);
true
}
fn bar_high_first(&self, bar: &Bar) -> bool {
!self.config.bar_adaptive_high_low_ordering || bar.high - bar.open < bar.open - bar.low
}
fn update_bar_quote_bid(&mut self, quote: &QuoteTick) {
let bid = BookOrder::new(
OrderSide::Buy,
quote.bid_price,
quote.bid_size,
OrderSide::Buy as u64,
);
self.book
.add(bid, 0, self.book.sequence.saturating_add(1), quote.ts_event);
}
fn clear_bar_quote_bid(&mut self, quote: &QuoteTick) {
self.book
.clear_bids(self.book.sequence.saturating_add(1), quote.ts_event);
}
fn update_bar_quote_ask(&mut self, quote: &QuoteTick) {
let ask = BookOrder::new(
OrderSide::Sell,
quote.ask_price,
quote.ask_size,
OrderSide::Sell as u64,
);
self.book
.add(ask, 0, self.book.sequence.saturating_add(1), quote.ts_event);
}
fn clear_bar_quote_ask(&mut self, quote: &QuoteTick) {
self.book
.clear_asks(self.book.sequence.saturating_add(1), quote.ts_event);
}
pub fn process_trade_tick(&mut self, trade: &TradeTick) {
log::debug!("Processing {trade}");
if let Err(e) = self.check_price_precision(trade.price.precision, "trade price") {
self.log_precision_mismatch("trade tick", trade.instrument_id, &e);
return;
}
if let Err(e) = self.check_size_precision(trade.size.precision, "trade size") {
self.log_precision_mismatch("trade tick", trade.instrument_id, &e);
return;
}
self.precision_mismatch_streak = 0;
let price_raw = trade.price.raw();
if self.book_type == BookType::L1_MBP {
if trade.ts_event < self.book.ts_last {
log::warn!(
"Skipping stale trade: ts_event {} < book.ts_last {} for {}",
trade.ts_event,
self.book.ts_last,
self.book.instrument_id,
);
self.iterate(trade.ts_init, AggressorSide::NoAggressor);
return;
}
if !self.update_trade_tick_or_skip(trade, "trade tick") {
return;
}
}
self.core.set_last_raw(trade.price);
if !self.config.trade_execution {
if self.book_type == BookType::L1_MBP {
if let Some(bid) = self.book.best_bid_price() {
self.core.set_bid_raw(bid);
}
if let Some(ask) = self.book.best_ask_price() {
self.core.set_ask_raw(ask);
}
} else {
self.iterate_with_mode(
trade.ts_init,
AggressorSide::NoAggressor,
OrderMatchMode::LastPriceStopTriggers,
);
}
return;
}
let aggressor_side = trade.aggressor_side;
match aggressor_side {
AggressorSide::Buy => {
if self.core.ask.is_none_or(|ask| trade.price > ask) {
self.core.set_ask_raw(trade.price);
}
if self.core.bid.is_none() {
self.core.set_bid_raw(trade.price);
}
}
AggressorSide::Sell => {
if self.core.bid.is_none_or(|bid| trade.price < bid) {
self.core.set_bid_raw(trade.price);
}
if self.core.ask.is_none() {
self.core.set_ask_raw(trade.price);
}
}
AggressorSide::NoAggressor => {
if self.core.bid.is_none_or(|bid| trade.price <= bid) {
self.core.set_bid_raw(trade.price);
}
if self.core.ask.is_none_or(|ask| trade.price >= ask) {
self.core.set_ask_raw(trade.price);
}
}
}
let original_bid = self.core.bid;
let original_ask = self.core.ask;
match aggressor_side {
AggressorSide::Sell => {
if original_ask.is_some_and(|ask| trade.price < ask) {
self.core.set_ask_raw(trade.price);
}
}
AggressorSide::Buy => {
if original_bid.is_some_and(|bid| trade.price > bid) {
self.core.set_bid_raw(trade.price);
}
}
AggressorSide::NoAggressor => {
self.core.set_bid_raw(trade.price);
self.core.set_ask_raw(trade.price);
}
}
self.last_trade_size = Some(trade.size);
self.trade_consumption = 0;
if self.config.liquidity_consumption && self.book_type != BookType::L1_MBP {
self.seed_trade_consumption(
price_raw,
trade.size.raw(),
trade.ts_event,
aggressor_side,
);
}
self.resolve_pending_on_trade(price_raw);
self.decrement_queue_on_trade(price_raw, trade.size.raw(), aggressor_side);
self.iterate(trade.ts_init, aggressor_side);
self.last_trade_size = None;
self.trade_consumption = 0;
if self.book_type == BookType::L1_MBP {
match aggressor_side {
AggressorSide::Sell => {
if let Some(ask) = self.last_quote_ask {
self.core.ask = Some(ask);
}
}
AggressorSide::Buy => {
if let Some(bid) = self.last_quote_bid {
self.core.bid = Some(bid);
}
}
AggressorSide::NoAggressor => {}
}
} else {
match aggressor_side {
AggressorSide::Sell => {
if let Some(ask) = original_ask
&& trade.price < ask
{
self.core.ask = Some(ask);
}
}
AggressorSide::Buy => {
if let Some(bid) = original_bid
&& trade.price > bid
{
self.core.bid = Some(bid);
}
}
AggressorSide::NoAggressor => {}
}
}
}
fn update_quote_tick_or_skip(&mut self, quote: &QuoteTick, context: &str) -> bool {
if let Err(e) = self.book.update_quote_tick(quote) {
log::warn!(
"Skipping {context} for {}: update_quote_tick failed: {e}",
quote.instrument_id,
);
return false;
}
true
}
fn update_trade_tick_or_skip(&mut self, trade: &TradeTick, context: &str) -> bool {
if let Err(e) = self.book.update_trade_tick(trade) {
log::warn!(
"Skipping {context} for {}: update_trade_tick failed: {e}",
trade.instrument_id,
);
return false;
}
true
}
pub fn process_status(&mut self, action: MarketStatusAction) {
log::debug!("Processing {action}");
match action {
MarketStatusAction::Trading | MarketStatusAction::PreOpen
if matches!(
self.market_status,
MarketStatus::Closed | MarketStatus::Paused | MarketStatus::Suspended
) =>
{
self.market_status = MarketStatus::Open;
}
MarketStatusAction::Pause if self.market_status == MarketStatus::Open => {
self.market_status = MarketStatus::Paused;
}
MarketStatusAction::Suspend if self.market_status == MarketStatus::Open => {
self.market_status = MarketStatus::Suspended;
}
MarketStatusAction::Halt | MarketStatusAction::Close
if self.market_status == MarketStatus::Open =>
{
self.market_status = MarketStatus::Closed;
}
_ => {}
}
}
pub fn process_instrument_close(&mut self, close: InstrumentClose) {
if close.instrument_id != self.instrument.id() {
log::warn!(
"Received instrument close for unknown instrument_id: {}",
close.instrument_id
);
return;
}
if close.close_type == InstrumentCloseType::ContractExpired {
self.instrument_close = Some(close);
self.iterate(close.ts_init, AggressorSide::NoAggressor);
}
}
pub fn process_instrument_expiration(&mut self, timestamp_ns: UnixNanos) {
self.check_instrument_expiration(timestamp_ns, false);
}
#[must_use]
pub const fn is_expiration_processed(&self) -> bool {
self.expiration_processed
}
fn requires_pending_resolution(&self) -> bool {
matches!(self.instrument, InstrumentAny::BinaryOption(_))
}
fn cancel_open_orders_for_expiration(&mut self) {
let instrument_id = self.instrument.id();
let expiration_order_ids: IndexSet<ClientOrderId> = {
let cache = self.cache.borrow();
let mut order_ids = IndexSet::new();
for order_info in self.get_open_orders() {
order_ids.insert(order_info.client_order_id);
}
for order in cache.orders(None, Some(&instrument_id), None, None, None) {
if order.is_open() || order.is_inflight() {
order_ids.insert(order.client_order_id());
}
}
order_ids
};
for client_order_id in expiration_order_ids {
let order = {
let cache = self.cache.borrow();
cache.order(&client_order_id).map(|order| order.clone())
};
if let Some(order) = order {
self.cancel_order(&order, None);
}
}
}
fn enter_pending_resolution(&mut self) {
if self.pending_resolution {
return;
}
self.pending_resolution = true;
self.market_status = MarketStatus::Closed;
self.cancel_open_orders_for_expiration();
log::info!(
"{} expired and is now pending resolution; open orders canceled and new orders blocked",
self.instrument.id()
);
}
fn check_instrument_expiration(&mut self, timestamp_ns: UnixNanos, defer_settlement: bool) {
if self.expiration_processed || self.option_settlement_failed {
return;
}
let timestamp_triggered = self
.instrument
.expiration_ns()
.is_some_and(|ns| timestamp_ns >= ns);
if !timestamp_triggered && self.instrument_close.is_none() {
return;
}
if self.instrument_close.is_none()
&& timestamp_triggered
&& self.requires_pending_resolution()
{
self.enter_pending_resolution();
return;
}
if matches!(
self.instrument,
InstrumentAny::OptionContract(_) | InstrumentAny::CryptoOption(_)
) {
if !self.option_expiration_orders_canceled {
self.option_expiration_orders_canceled = true;
self.enter_pending_resolution();
}
if defer_settlement
&& self.instrument_close.is_none()
&& self.instrument.expiration_ns() == Some(timestamp_ns)
{
return;
}
match self.process_option_expiry(timestamp_ns) {
Ok(true) => {
self.expiration_processed = true;
self.pending_resolution = false;
self.instrument_close.take();
self.option_settlement_warning = None;
log::info!("{} reached expiration", self.instrument.id());
}
Ok(false) => {}
Err(e) => {
self.option_settlement_failed = true;
log::error!(
"Option settlement failed terminally for {}: {e}",
self.instrument.id()
);
}
}
return;
}
self.expiration_processed = true;
self.pending_resolution = false;
let close = self.instrument_close.take();
log::info!("{} reached expiration", self.instrument.id());
self.cancel_open_orders_for_expiration();
let instrument_id = self.instrument.id();
let positions: Vec<(
TraderId,
StrategyId,
AccountId,
PositionId,
OrderSide,
Quantity,
)> = {
let cache = self.cache.borrow();
cache
.positions_open(None, Some(&instrument_id), None, None, None)
.into_iter()
.filter_map(|pos| {
OrderCore::closing_side(pos.side).map(|closing_side| {
(
pos.trader_id,
pos.strategy_id,
pos.account_id,
pos.id,
closing_side,
pos.quantity,
)
})
})
.collect()
};
let ts_now = self.clock.borrow().timestamp_ns();
let close_price = close.as_ref().map(|close| close.close_price);
for (trader_id, strategy_id, account_id, position_id, closing_side, quantity) in positions {
let client_order_id =
ClientOrderId::from(format!("EXPIRATION-{}-{}", self.venue, UUID4::new()).as_str());
let mut order = OrderAny::Market(MarketOrder::new(
trader_id,
strategy_id,
instrument_id,
client_order_id,
closing_side,
quantity,
TimeInForce::Gtc,
UUID4::new(),
ts_now,
true, false,
None,
None,
None,
None,
None,
None,
None,
Some(vec![Ustr::from(&format!(
"EXPIRATION_{}_CLOSE",
self.venue
))]),
));
order.set_liquidity_side(LiquiditySide::Taker);
let add_result =
self.cache
.borrow_mut()
.add_order(order.clone(), Some(position_id), None, false);
if add_result.is_err() {
log::debug!("Expiration order already in cache: {client_order_id}");
} else {
self.publish_order_initialized(&order);
}
let venue_order_id = self.ids_generator.get_venue_order_id(&order).unwrap();
self.account_ids.insert(trader_id, account_id);
self.generate_order_accepted(&order, venue_order_id);
if let Some(fill_price) = close_price {
if let Err(e) = self.apply_fills(
&order,
&[(fill_price, quantity)],
LiquiditySide::Taker,
Some(position_id),
None,
None,
) {
log::error!("Cannot fill expiration order {client_order_id}: {e}");
}
} else {
self.fill_market_order(client_order_id);
}
}
}
pub fn liquidate_open_positions(
&mut self,
ts_now: UnixNanos,
cancel_open_orders: bool,
settlement_currency: Currency,
) {
if self.instrument.settlement_currency() != settlement_currency {
return;
}
if cancel_open_orders {
let open_orders: Vec<RestingOrder> = self.get_open_orders();
for order_info in &open_orders {
let order = {
let cache = self.cache.borrow();
cache.order_owned(&order_info.client_order_id)
};
if let Some(order) = order {
self.cancel_order(&order, None);
}
}
}
let instrument_id = self.instrument.id();
let positions: Vec<(
TraderId,
StrategyId,
AccountId,
PositionId,
OrderSide,
Quantity,
)> = {
let cache = self.cache.borrow();
cache
.positions_open(None, Some(&instrument_id), None, None, None)
.into_iter()
.filter_map(|pos| {
OrderCore::closing_side(pos.side).map(|closing_side| {
(
pos.trader_id,
pos.strategy_id,
pos.account_id,
pos.id,
closing_side,
pos.quantity,
)
})
})
.collect()
};
for (trader_id, strategy_id, account_id, position_id, closing_side, quantity) in positions {
let has_price = if closing_side == OrderSide::Sell {
self.best_bid_price().is_some()
} else {
self.best_ask_price().is_some()
};
if !has_price {
log::warn!(
"LIQUIDATION: no price available for {instrument_id} position {position_id}, skipping"
);
continue;
}
let client_order_id = ClientOrderId::from(
format!("LIQUIDATION-{}-{}", self.venue, UUID4::new()).as_str(),
);
let order = OrderAny::Market(MarketOrder::new(
trader_id,
strategy_id,
instrument_id,
client_order_id,
closing_side,
quantity,
TimeInForce::Ioc,
UUID4::new(),
ts_now,
true, false,
None,
None,
None,
None,
None,
None,
None,
Some(vec![Ustr::from(&format!(
"LIQUIDATION_{}_CLOSE",
self.venue
))]),
));
let venue_order_id = self.ids_generator.get_venue_order_id(&order).unwrap();
{
let mut cache = self.cache.borrow_mut();
if let Err(e) = cache.add_order(order.clone(), Some(position_id), None, false) {
log::debug!("Liquidation order already in cache: {e}");
} else {
drop(cache);
self.publish_order_initialized(&order);
self.cache
.borrow_mut()
.add_venue_order_id(&client_order_id, &venue_order_id, false)
.ok();
}
}
self.account_ids.insert(trader_id, account_id);
self.generate_order_submitted(&order, account_id);
self.generate_order_accepted(&order, venue_order_id);
self.fill_market_order(client_order_id);
}
}
pub fn process_order(&mut self, order: &mut OrderAny, account_id: AccountId) {
if self.core.order_exists(order.client_order_id()) {
return;
}
let ts_now = self.clock.borrow().timestamp_ns();
self.check_instrument_expiration(ts_now, self.config.defer_option_settlement);
let reject_reason: Option<Ustr> = 'validate: {
let cache_borrow = self.cache.as_ref().borrow();
self.account_ids.insert(order.trader_id(), account_id);
if self.pending_resolution {
break 'validate Some(
format!(
"Contract {} has expired and is pending resolution",
self.instrument.id()
)
.into(),
);
}
if self.market_status != MarketStatus::Open {
break 'validate Some(
format!(
"Market {} is {}, cannot accept order {}",
self.instrument.id(),
self.market_status,
order.client_order_id()
)
.into(),
);
}
if self.instrument.has_expiration() {
if let Some(activation_ns) = self.instrument.activation_ns()
&& self.clock.borrow().timestamp_ns() < activation_ns
{
break 'validate Some(
format!(
"Contract {} is not yet active, activation {activation_ns}",
self.instrument.id(),
)
.into(),
);
}
if let Some(expiration_ns) = self.instrument.expiration_ns()
&& self.clock.borrow().timestamp_ns() >= expiration_ns
{
break 'validate Some(
format!(
"Contract {} has expired, expiration {expiration_ns}",
self.instrument.id(),
)
.into(),
);
}
}
if self.config.support_contingent_orders {
if let Some(parent_order_id) = order.parent_order_id() {
let parent_order = match self.order_snapshot(parent_order_id) {
Some(o) if o.contingency_type() == Some(ContingencyType::Oto) => o,
_ => panic!("OTO parent not found"),
};
let parent_filled_qty = parent_order.filled_qty();
if parent_order.status() == OrderStatus::Rejected && order.is_open() {
break 'validate Some(
format!("Rejected OTO order from {parent_order_id}").into(),
);
} else if parent_filled_qty.is_zero()
|| (self.config.oto_full_trigger
&& parent_filled_qty < parent_order.quantity())
{
log::info!(
"Pending OTO order {} triggers from {parent_order_id}",
order.client_order_id(),
);
return;
}
}
if let Some(linked_order_ids) = order.linked_order_ids() {
let contingency_type = order.contingency_type();
for client_order_id in linked_order_ids {
match cache_borrow.order(client_order_id) {
Some(contingent_order)
if matches!(
contingency_type,
Some(ContingencyType::Oco | ContingencyType::Ouo)
) && !order.is_closed()
&& contingent_order.is_closed() =>
{
break 'validate Some(
format!("Contingent order {client_order_id} already closed")
.into(),
);
}
None => panic!("Cannot find contingent order for {client_order_id}"),
_ => {}
}
}
}
}
if order.quantity().precision != self.instrument.size_precision() {
break 'validate Some(
format!(
"Invalid order quantity precision for order {}, was {} when {} size precision is {}",
order.client_order_id(),
order.quantity().precision,
self.instrument.id(),
self.instrument.size_precision()
)
.into(),
);
}
if let Some(display_qty) = order.display_qty()
&& display_qty.precision != self.instrument.size_precision()
{
break 'validate Some(
format!(
"Invalid order display quantity precision for order {}, was {} when {} size precision is {}",
order.client_order_id(),
display_qty.precision,
self.instrument.id(),
self.instrument.size_precision()
)
.into(),
);
}
if let Some(price) = order.price()
&& price.precision != self.instrument.price_precision()
{
break 'validate Some(
format!(
"Invalid order price precision for order {}, was {} when {} price precision is {}",
order.client_order_id(),
price.precision,
self.instrument.id(),
self.instrument.price_precision()
)
.into(),
);
}
if let Some(trigger_price) = order.trigger_price()
&& trigger_price.precision != self.instrument.price_precision()
{
break 'validate Some(
format!(
"Invalid order trigger price precision for order {}, was {} when {} price precision is {}",
order.client_order_id(),
trigger_price.precision,
self.instrument.id(),
self.instrument.price_precision()
)
.into(),
);
}
if order.is_reduce_only() && !self.config.use_reduce_only {
break 'validate Some(
"Reduce-only orders are not supported by this matching engine".into(),
);
}
let position = self.position_for_order_in_cache(&cache_borrow, order);
if order.order_side() == OrderSide::Sell
&& self.account_type != AccountType::Margin
&& matches!(self.instrument, InstrumentAny::Equity(_))
&& position
.as_ref()
.is_none_or(|pos| !order.would_reduce_only(pos.side, pos.quantity))
{
let position_string = position
.as_ref()
.map_or("None".to_string(), |pos| pos.id.to_string());
break 'validate Some(
format!(
"Short selling not permitted on a CASH account with position {position_string} and order {order}",
)
.into(),
);
}
if self.config.use_reduce_only
&& order.is_reduce_only()
&& !order.is_closed()
&& position.as_ref().is_none_or(|pos| {
pos.is_closed()
|| (order.is_buy() && pos.is_long())
|| (order.is_sell() && pos.is_short())
})
{
break 'validate Some(
format!(
"Reduce-only order {} ({}-{}) would have increased position",
order.client_order_id(),
order.order_type().to_string().to_uppercase(),
order.order_side().to_string().to_uppercase()
)
.into(),
);
}
None
};
if let Some(reason) = reject_reason {
self.generate_order_rejected(order, reason);
return;
}
if order.is_quote_quantity()
&& !self.instrument.is_inverse()
&& !matches!(
order.order_type(),
OrderType::TrailingStopLimit | OrderType::TrailingStopMarket,
)
&& (order.price().is_some()
|| matches!(
order.order_type(),
OrderType::Market | OrderType::MarketToLimit,
))
&& !self.convert_quote_to_base_quantity(order)
{
return;
}
match order.order_type() {
OrderType::Market => self.process_market_order(order),
OrderType::Limit => self.process_limit_order(order),
OrderType::MarketToLimit => self.process_market_to_limit_order(order),
OrderType::StopMarket => self.process_stop_market_order(order),
OrderType::StopLimit => self.process_stop_limit_order(order),
OrderType::MarketIfTouched => self.process_market_if_touched_order(order),
OrderType::LimitIfTouched => self.process_limit_if_touched_order(order),
OrderType::TrailingStopMarket => self.process_trailing_stop_order(order),
OrderType::TrailingStopLimit => self.process_trailing_stop_order(order),
}
}
fn convert_quote_to_base_quantity(&self, order: &mut OrderAny) -> bool {
let reference_price = if let Some(price) = order.price() {
Some(price)
} else {
match order.order_side() {
OrderSide::Buy => self.core.ask,
OrderSide::Sell => self.core.bid,
}
};
let Some(reference_price) = reference_price else {
self.generate_order_rejected(
order,
format!(
"No market for {} to convert quote quantity to base",
order.instrument_id(),
)
.into(),
);
return false;
};
let base_quantity = self
.instrument
.calculate_base_quantity(order.quantity(), reference_price);
let ts_now = self.clock.borrow().timestamp_ns();
let event = OrderEventAny::Updated(OrderUpdated::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
base_quantity,
UUID4::new(),
ts_now,
ts_now,
false,
order.venue_order_id(),
order.account_id(),
None,
None,
None,
false,
));
if let Err(e) = order.apply(event.clone()) {
log::error!(
"Failed to apply quote-to-base update for {}: {e}",
order.client_order_id(),
);
return false;
}
self.dispatch_order_event(event);
true
}
pub fn process_modify(&mut self, command: &ModifyOrder, account_id: AccountId) {
if !self.core.order_exists(command.client_order_id) {
self.generate_order_modify_rejected(
command.trader_id,
command.strategy_id,
command.instrument_id,
command.client_order_id,
Ustr::from(format!("Order {} not found", command.client_order_id).as_str()),
command.venue_order_id,
Some(account_id),
);
return;
}
let order = match self.order_snapshot(command.client_order_id) {
Some(order) => order,
None => {
log::error!(
"Cannot modify order: order {} not found in cache",
command.client_order_id
);
return;
}
};
let update_success = self.update_order(
&order,
command.quantity,
command.price,
command.trigger_price,
None,
);
if !update_success {
return;
}
if !self.core.order_exists(command.client_order_id) {
return;
}
let Some(refreshed) = self.resync_core_entry(command.client_order_id) else {
return;
};
let price_changed = refreshed.price() != order.price()
|| refreshed.trigger_price() != order.trigger_price();
if price_changed
&& refreshed.is_open()
&& self.config.queue_position
&& let Some(new_price) = refreshed.price()
{
self.snapshot_queue_position(&refreshed, new_price);
self.queue_excess.swap_remove(&refreshed.client_order_id());
}
}
pub fn process_cancel(&mut self, command: &CancelOrder, account_id: AccountId) {
if !self.core.order_exists(command.client_order_id) {
self.generate_order_cancel_rejected(
command.trader_id,
command.strategy_id,
account_id,
command.instrument_id,
command.client_order_id,
command.venue_order_id,
Ustr::from(format!("Order {} not found", command.client_order_id).as_str()),
);
return;
}
let order = match self.order_snapshot(command.client_order_id) {
Some(order) => order,
None => {
log::error!(
"Cannot cancel order: order {} not found in cache",
command.client_order_id
);
return;
}
};
if !order.is_inflight() && !order.is_open() {
self.purge_stale_core_entry(command.client_order_id);
return;
}
self.cancel_order(&order, None);
}
pub fn process_cancel_all(&mut self, command: &CancelAllOrders, account_id: AccountId) {
self.process_cancel_all_excluding(command, account_id, &[]);
}
pub fn process_cancel_all_excluding(
&mut self,
command: &CancelAllOrders,
account_id: AccountId,
excluded: &[ClientOrderId],
) {
let instrument_id = command.instrument_id;
let order_side = command.order_side;
let mut client_order_ids: Vec<ClientOrderId> = {
let cache = self.cache.borrow();
cache
.orders_open_refs(
None,
Some(&instrument_id),
None,
Some(&account_id),
order_side,
)
.into_iter()
.chain(cache.orders_inflight_refs(
None,
Some(&instrument_id),
None,
Some(&account_id),
order_side,
))
.map(|order| order.client_order_id())
.filter(|client_order_id| !excluded.contains(client_order_id))
.collect()
};
client_order_ids.sort_unstable();
client_order_ids.dedup();
for client_order_id in client_order_ids {
let order = match self
.cache
.borrow()
.order(&client_order_id)
.map(|o| o.clone())
{
Some(order) => order,
None => continue,
};
if !order.is_inflight() && !order.is_open() {
self.purge_stale_core_entry(client_order_id);
continue;
}
self.cancel_order_excluding(&order, None, excluded);
}
}
fn purge_stale_core_entry(&mut self, client_order_id: ClientOrderId) {
if self.core.order_exists(client_order_id) {
self.delete_core_order(client_order_id);
}
self.remove_queue_position(client_order_id);
self.cached_filled_qty.swap_remove(&client_order_id);
}
fn resync_core_entry(&mut self, client_order_id: ClientOrderId) -> Option<OrderAny> {
let order = self.order_snapshot(client_order_id)?;
if order.is_closed() {
self.delete_core_order(client_order_id);
self.remove_queue_position(client_order_id);
return Some(order);
}
let new_match_info = Self::matching_core_entry(&order);
let unchanged = self
.core
.get_order(client_order_id)
.is_some_and(|existing| *existing == new_match_info);
if unchanged {
self.track_post_match_order(&order);
return Some(order);
}
self.delete_core_order(client_order_id);
self.track_post_match_order(&order);
self.core.add_order(new_match_info);
Some(order)
}
fn order_snapshot(&self, client_order_id: ClientOrderId) -> Option<OrderAny> {
let mut order = self.cache.borrow().order(&client_order_id)?.clone();
let mut pending = self.pending_order_updates.borrow_mut();
if order.is_closed() {
pending.swap_remove(&client_order_id);
return Some(order);
}
if let Some(updates) = pending.get_mut(&client_order_id) {
Self::retain_unapplied_order_updates(&order, updates);
for update in updates.iter() {
if let Err(e) = order.apply(OrderEventAny::Updated(*update)) {
log::error!("Cannot apply pending update for {client_order_id}: {e}");
return None;
}
}
if updates.is_empty() {
pending.swap_remove(&client_order_id);
}
}
if let Some(filled_qty) = self.cached_filled_qty.get(&client_order_id) {
write_filled_qty(&mut order, *filled_qty);
order.set_leaves_qty(order.quantity().saturating_sub(*filled_qty));
}
Some(order)
}
fn purge_applied_order_updates(&self) {
let cache = self.cache.borrow();
self.pending_order_updates
.borrow_mut()
.retain(|id, updates| {
let Some(order) = cache.order(id) else {
return false;
};
Self::retain_unapplied_order_updates(&order, updates);
!updates.is_empty()
});
}
fn retain_unapplied_order_updates(order: &OrderAny, updates: &mut Vec<OrderUpdated>) {
if order.is_closed() {
updates.clear();
return;
}
let events = order.events();
updates.retain(|update| {
!events.iter().any(|event| {
matches!(event, OrderEventAny::Updated(applied) if applied.event_id == update.event_id)
})
});
}
pub fn process_batch_cancel(&mut self, command: &BatchCancelOrders, account_id: AccountId) {
for order in &command.cancels {
self.process_cancel(order, account_id);
}
}
pub fn process_batch_modify(&mut self, command: &BatchModifyOrders, account_id: AccountId) {
for order in &command.modifies {
self.process_modify(order, account_id);
}
}
fn process_market_order(&mut self, order: &OrderAny) {
if order.time_in_force() == TimeInForce::AtTheOpen
|| order.time_in_force() == TimeInForce::AtTheClose
{
self.generate_order_rejected(
order,
format!(
"time in force {} is not currently supported",
order.time_in_force()
)
.into(),
);
return;
}
if (order.order_side() == OrderSide::Buy && self.core.ask.is_none())
|| (order.order_side() == OrderSide::Sell && self.core.bid.is_none())
{
self.generate_order_rejected(
order,
format!("No market for {}", order.instrument_id()).into(),
);
return;
}
if self.config.use_market_order_acks {
let venue_order_id = self.ids_generator.get_venue_order_id(order).unwrap();
self.generate_order_accepted(order, venue_order_id);
}
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
self.fill_market_order(order.client_order_id());
}
fn process_limit_order(&mut self, order: &mut OrderAny) {
if order.time_in_force() == TimeInForce::AtTheOpen
|| order.time_in_force() == TimeInForce::AtTheClose
{
self.generate_order_rejected(
order,
format!(
"time in force {} is not currently supported",
order.time_in_force()
)
.into(),
);
return;
}
let limit_px = order.price().expect("Limit order must have a price");
if order.is_post_only() && self.core.is_limit_matched(order.order_side(), limit_px) {
self.generate_order_rejected(
order,
format!(
"POST_ONLY {} {} order limit px of {} would have been a TAKER: bid={}, ask={}",
order.order_type(),
order.order_side(),
order.price().unwrap(),
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
)
.into(),
);
return;
}
self.accept_order(order);
if self.core.is_limit_matched(order.order_side(), limit_px) {
order.set_liquidity_side(LiquiditySide::Taker);
if self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
.is_err()
&& let Err(e) = self.cache.borrow_mut().replace_order(order)
{
log::debug!("Failed to update order in cache: {e}");
}
self.fill_limit_order(order.client_order_id());
if self.core.order_exists(order.client_order_id())
&& let Some(mut order) = self.cache.borrow_mut().order_mut(&order.client_order_id())
{
order.set_liquidity_side(LiquiditySide::Maker);
}
} else if matches!(order.time_in_force(), TimeInForce::Fok | TimeInForce::Ioc) {
self.cancel_order(order, None);
} else {
order.set_liquidity_side(LiquiditySide::Maker);
if let Some(price) = order.price() {
self.snapshot_queue_position(order, price);
}
let add_result = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false);
if let Err(e) = add_result {
log::debug!("Failed to add order to cache: {e}");
if let Some(mut order) = self.cache.borrow_mut().order_mut(&order.client_order_id())
&& !matches!(
order.liquidity_side(),
Some(LiquiditySide::Maker | LiquiditySide::Taker)
)
{
order.set_liquidity_side(LiquiditySide::Maker);
}
}
}
}
fn process_market_to_limit_order(&mut self, order: &OrderAny) {
if (order.order_side() == OrderSide::Buy && self.core.ask.is_none())
|| (order.order_side() == OrderSide::Sell && self.core.bid.is_none())
{
self.generate_order_rejected(
order,
format!("No market for {}", order.instrument_id()).into(),
);
return;
}
if self.config.use_market_order_acks {
let venue_order_id = self.ids_generator.get_venue_order_id(order).unwrap();
self.generate_order_accepted(order, venue_order_id);
}
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
let client_order_id = order.client_order_id();
self.fill_market_order(client_order_id);
let filled_qty = self
.cached_filled_qty
.get(&client_order_id)
.copied()
.unwrap_or_default();
let leaves_qty = order.quantity().saturating_sub(filled_qty);
if leaves_qty.is_zero() {
self.purge_cached_filled_qty_if_closed(client_order_id);
return;
}
if let Some(mut updated_order) = self.order_snapshot(client_order_id) {
self.accept_order(&mut updated_order);
}
}
fn process_stop_market_order(&mut self, order: &mut OrderAny) {
let stop_px = order
.trigger_price()
.expect("Stop order must have a trigger price");
if self.core.is_stop_matched_with_trigger_type(
order.order_side(),
stop_px,
order.trigger_type().unwrap_or(TriggerType::Default),
) {
if self.config.reject_stop_orders {
self.generate_order_rejected(
order,
format!(
"{} {} order stop px of {} was in the market: bid={}, ask={}, but rejected because of configuration",
order.order_type(),
order.order_side(),
order.trigger_price().unwrap(),
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
).into(),
);
return;
}
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
self.fill_market_order(order.client_order_id());
return;
}
self.accept_order(order);
order.set_liquidity_side(LiquiditySide::Maker);
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
}
fn process_stop_limit_order(&mut self, order: &mut OrderAny) {
let stop_px = order
.trigger_price()
.expect("Stop order must have a trigger price");
if self.core.is_stop_matched_with_trigger_type(
order.order_side(),
stop_px,
order.trigger_type().unwrap_or(TriggerType::Default),
) {
if self.config.reject_stop_orders {
self.generate_order_rejected(
order,
format!(
"{} {} order stop px of {} was in the market: bid={}, ask={}, but rejected because of configuration",
order.order_type(),
order.order_side(),
order.trigger_price().unwrap(),
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
).into(),
);
return;
}
self.accept_triggered_limit_style_order(order);
return;
}
self.accept_order(order);
order.set_liquidity_side(LiquiditySide::Maker);
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
}
fn process_market_if_touched_order(&mut self, order: &mut OrderAny) {
if self.core.is_touch_triggered_with_trigger_type(
order.order_side(),
order.trigger_price().unwrap(),
order.trigger_type().unwrap_or(TriggerType::Default),
) {
if self.config.reject_stop_orders {
self.generate_order_rejected(
order,
format!(
"{} {} order trigger px of {} was in the market: bid={}, ask={}, but rejected because of configuration",
order.order_type(),
order.order_side(),
order.trigger_price().unwrap(),
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
).into(),
);
return;
}
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
self.fill_market_order(order.client_order_id());
return;
}
self.accept_order(order);
order.set_liquidity_side(LiquiditySide::Maker);
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
}
fn process_limit_if_touched_order(&mut self, order: &mut OrderAny) {
if self.core.is_touch_triggered_with_trigger_type(
order.order_side(),
order.trigger_price().unwrap(),
order.trigger_type().unwrap_or(TriggerType::Default),
) {
if self.config.reject_stop_orders {
self.generate_order_rejected(
order,
format!(
"{} {} order trigger px of {} was in the market: bid={}, ask={}, but rejected because of configuration",
order.order_type(),
order.order_side(),
order.trigger_price().unwrap(),
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
).into(),
);
return;
}
self.accept_triggered_limit_style_order(order);
return;
}
self.accept_order(order);
order.set_liquidity_side(LiquiditySide::Maker);
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
}
fn accept_triggered_limit_style_order(&mut self, order: &mut OrderAny) {
self.accept_order(order);
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
self.trigger_limit_style_stop_order(order.client_order_id(), order.clone());
if let Some(cached_order) = self
.cache
.borrow()
.order(&order.client_order_id())
.map(|order| order.clone())
{
*order = cached_order;
}
}
fn process_trailing_stop_order(&mut self, order: &mut OrderAny) {
if let Some(trigger_price) = order.trigger_price()
&& self.core.is_stop_matched_with_trigger_type(
order.order_side(),
trigger_price,
order.trigger_type().unwrap_or(TriggerType::Default),
)
{
self.generate_order_rejected(
order,
format!(
"{} {} order trigger px of {} was in the market: bid={}, ask={}, but rejected because of configuration",
order.order_type(),
order.order_side(),
trigger_price,
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
).into(),
);
return;
}
order.set_liquidity_side(LiquiditySide::Maker);
self.accept_order(order);
if let Err(e) = self
.cache
.borrow_mut()
.add_order(order.clone(), None, None, false)
{
log::debug!("Order already in cache: {e}");
}
}
pub fn iterate(&mut self, timestamp_ns: UnixNanos, aggressor_side: AggressorSide) {
self.iterate_with_mode(timestamp_ns, aggressor_side, OrderMatchMode::All);
}
fn iterate_with_mode(
&mut self,
timestamp_ns: UnixNanos,
aggressor_side: AggressorSide,
match_mode: OrderMatchMode,
) {
self.purge_closed_cached_filled_qty();
self.purge_applied_order_updates();
self.purge_applied_fills();
if aggressor_side == AggressorSide::NoAggressor && self.last_trade_size.is_none() {
if self.book_type == BookType::L1_MBP {
if let Some(bid) = self.book.best_bid_price() {
self.core.set_bid_raw(bid);
}
if let Some(ask) = self.book.best_ask_price() {
self.core.set_ask_raw(ask);
}
} else {
self.core.bid = self.book.best_bid_price();
self.core.ask = self.book.best_ask_price();
}
}
let mut matched_order = false;
if self.market_status == MarketStatus::Open {
for action in self.core.iterate_bids() {
if !self.should_process_match_action(action, match_mode) {
continue;
}
matched_order = true;
match action {
MatchAction::FillLimit(id) => self.fill_resting_limit_order(id),
MatchAction::TriggerStop(id) => self.trigger_stop_order(id),
}
}
for action in self.core.iterate_asks() {
if !self.should_process_match_action(action, match_mode) {
continue;
}
matched_order = true;
match action {
MatchAction::FillLimit(id) => self.fill_resting_limit_order(id),
MatchAction::TriggerStop(id) => self.trigger_stop_order(id),
}
}
}
let order_ids: Vec<ClientOrderId> = if matched_order {
self.core.iter_orders().map(|m| m.client_order_id).collect()
} else if self.post_match_order_ids.is_empty() {
Vec::new()
} else {
self.core
.iter_orders()
.filter_map(|order| {
self.post_match_order_ids
.contains(&order.client_order_id)
.then_some(order.client_order_id)
})
.collect()
};
let support_gtd_orders = self.config.support_gtd_orders;
for client_order_id in order_ids {
let (action, keep_tracking) = {
let cache = self.cache.borrow();
let Some(order) = cache.order(&client_order_id) else {
self.post_match_order_ids.swap_remove(&client_order_id);
continue;
};
(
post_match_order_action(&order, support_gtd_orders, timestamp_ns, |order| {
self.order_snapshot(client_order_id)
.unwrap_or_else(|| order.clone())
}),
Self::requires_post_match_maintenance(&order),
)
};
match action {
PostMatchOrderAction::RemoveClosed => {
self.delete_core_order(client_order_id);
self.remove_queue_position(client_order_id);
self.cached_filled_qty.swap_remove(&client_order_id);
continue;
}
PostMatchOrderAction::Expire(order) => {
self.delete_core_order(client_order_id);
self.cached_filled_qty.swap_remove(&client_order_id);
self.expire_order(&order);
continue;
}
PostMatchOrderAction::UpdateTrailing(mut order) => {
if self.maybe_activate_trailing_stop(
&mut order,
self.core.bid,
self.core.ask,
self.core.last,
) {
self.update_trailing_stop_order(&order);
self.resync_core_entry(client_order_id);
}
}
PostMatchOrderAction::NoMaintenance => {
if !keep_tracking {
self.post_match_order_ids.swap_remove(&client_order_id);
}
}
}
if self.target_bid.is_some() || self.target_ask.is_some() || self.target_last.is_some()
{
if let Some(t) = self.target_bid.take() {
self.core.bid = Some(t);
}
if let Some(t) = self.target_ask.take() {
self.core.ask = Some(t);
}
if let Some(t) = self.target_last.take() {
self.core.last = Some(t);
}
}
}
if let Some(t) = self.target_bid.take() {
self.core.bid = Some(t);
}
if let Some(t) = self.target_ask.take() {
self.core.ask = Some(t);
}
if let Some(t) = self.target_last.take() {
self.core.last = Some(t);
}
self.core.bid = self.book.best_bid_price();
self.core.ask = self.book.best_ask_price();
self.check_instrument_expiration(timestamp_ns, self.config.defer_option_settlement);
self.purge_closed_cached_filled_qty();
self.purge_applied_order_updates();
self.purge_applied_fills();
}
fn fill_resting_limit_order(&mut self, client_order_id: ClientOrderId) {
if self
.core
.get_order(client_order_id)
.is_some_and(|order| order.order_type == OrderType::MarketToLimit)
&& let Some(mut order) = self.cache.borrow_mut().order_mut(&client_order_id)
{
order.set_liquidity_side(LiquiditySide::Maker);
}
self.fill_limit_order(client_order_id);
}
fn should_process_match_action(&self, action: MatchAction, match_mode: OrderMatchMode) -> bool {
let client_order_id = match action {
MatchAction::FillLimit(id) | MatchAction::TriggerStop(id) => id,
};
if !self.core.order_exists(client_order_id) {
return false;
}
match match_mode {
OrderMatchMode::All => true,
OrderMatchMode::LastPriceStopTriggers => match action {
MatchAction::TriggerStop(client_order_id) => self
.core
.get_order(client_order_id)
.is_some_and(|order| order.trigger_type == Some(TriggerType::LastPrice)),
MatchAction::FillLimit(_) => false,
},
}
}
fn get_trailing_activation_price(
&self,
trigger_type: TriggerType,
order_side: OrderSide,
bid: Option<Price>,
ask: Option<Price>,
last: Option<Price>,
) -> Option<Price> {
match trigger_type {
TriggerType::LastPrice => last,
TriggerType::LastOrBidAsk => last.or(match order_side {
OrderSide::Buy => ask,
OrderSide::Sell => bid,
}),
_ => match order_side {
OrderSide::Buy => ask,
OrderSide::Sell => bid,
},
}
}
fn maybe_activate_trailing_stop(
&self,
order: &mut OrderAny,
bid: Option<Price>,
ask: Option<Price>,
last: Option<Price>,
) -> bool {
match order {
OrderAny::TrailingStopMarket(inner) => {
if inner.is_activated {
return true;
}
if inner.activation_price.is_none() {
let px = self.get_trailing_activation_price(
inner.trigger_type,
inner.order_side(),
bid,
ask,
last,
);
if let Some(p) = px {
inner.activation_price = Some(p);
inner.set_activated();
if let Err(e) = self.cache.borrow_mut().replace_order(order) {
log::error!("Failed to update order: {e}");
}
return true;
}
return false;
}
let activation_price = inner.activation_price.unwrap();
let hit = match inner.order_side() {
OrderSide::Buy => ask.is_some_and(|a| a <= activation_price),
OrderSide::Sell => bid.is_some_and(|b| b >= activation_price),
};
if hit {
inner.set_activated();
if let Err(e) = self.cache.borrow_mut().replace_order(order) {
log::error!("Failed to update order: {e}");
}
}
hit
}
OrderAny::TrailingStopLimit(inner) => {
if inner.is_activated {
return true;
}
if inner.activation_price.is_none() {
let px = self.get_trailing_activation_price(
inner.trigger_type,
inner.order_side(),
bid,
ask,
last,
);
if let Some(p) = px {
inner.activation_price = Some(p);
inner.set_activated();
if let Err(e) = self.cache.borrow_mut().replace_order(order) {
log::error!("Failed to update order: {e}");
}
return true;
}
return false;
}
let activation_price = inner.activation_price.unwrap();
let hit = match inner.order_side() {
OrderSide::Buy => ask.is_some_and(|a| a <= activation_price),
OrderSide::Sell => bid.is_some_and(|b| b >= activation_price),
};
if hit {
inner.set_activated();
if let Err(e) = self.cache.borrow_mut().replace_order(order) {
log::error!("Failed to update order: {e}");
}
}
hit
}
_ => true,
}
}
fn determine_limit_price_and_volume(&mut self, order: &OrderAny) -> Vec<(Price, Quantity)> {
match order.price() {
Some(order_price) => {
let mut fills = if self.config.liquidity_consumption {
let size_prec = self.instrument.size_precision();
self.book
.get_all_crossed_levels(order.order_side(), order_price, size_prec)
} else {
let book_order =
BookOrder::new(order.order_side(), order_price, order.quantity(), 1);
self.book.simulate_fills(&book_order)
};
if let Some(trade_size) = self.last_trade_size
&& let Some(trade_price) = self.core.last
{
let fills_at_trade_price = fills.iter().any(|(px, _)| *px == trade_price);
if !fills_at_trade_price
&& self.core.is_limit_matched(order.order_side(), order_price)
{
let leaves_qty = order.leaves_qty();
let available_qty = if self.config.liquidity_consumption {
let remaining = trade_size.raw().saturating_sub(self.trade_consumption);
Quantity::from_raw(remaining, trade_size.precision)
} else {
trade_size
};
let fill_qty = min(leaves_qty, available_qty);
if fill_qty.non_zero() {
log::debug!(
"Trade execution fill: {} @ {} (trade_price={}, available: {}, book had {} fills)",
fill_qty,
order_price,
trade_price,
available_qty,
fills.len()
);
if self.config.liquidity_consumption {
self.trade_consumption += fill_qty.raw();
}
return vec![(order_price, fill_qty)];
}
}
}
if fills.is_empty() {
return fills;
}
let book_prices: Vec<Price> = if self.config.liquidity_consumption {
fills.iter().map(|(px, _)| *px).collect()
} else {
Vec::new()
};
let book_prices_ref: Option<&[Price]> = if book_prices.is_empty() {
None
} else {
Some(&book_prices)
};
if order
.liquidity_side()
.is_some_and(|liquidity_side| liquidity_side == LiquiditySide::Maker)
{
match order.order_side() {
OrderSide::Buy => {
let target_price = if order
.trigger_price()
.is_some_and(|trigger_price| order_price > trigger_price)
{
order.trigger_price().unwrap()
} else {
order_price
};
for fill in &mut fills {
let last_px = fill.0;
if last_px < order_price {
self.target_bid = self.core.bid;
self.target_ask = self.core.ask;
self.target_last = self.core.last;
self.core.set_ask_raw(target_price);
self.core.set_last_raw(target_price);
fill.0 = target_price;
}
}
}
OrderSide::Sell => {
let target_price = if order
.trigger_price()
.is_some_and(|trigger_price| order_price < trigger_price)
{
order.trigger_price().unwrap()
} else {
order_price
};
for fill in &mut fills {
let last_px = fill.0;
if last_px > order_price {
self.target_bid = self.core.bid;
self.target_ask = self.core.ask;
self.target_last = self.core.last;
self.core.set_bid_raw(target_price);
self.core.set_last_raw(target_price);
fill.0 = target_price;
}
}
}
}
}
self.apply_liquidity_consumption(
fills,
order.order_side(),
order.leaves_qty(),
book_prices_ref,
)
}
None => panic!("Limit order must have a price"),
}
}
fn determine_market_price_and_volume(&self, order: &OrderAny) -> Vec<(Price, Quantity)> {
let price = match order.order_side() {
OrderSide::Buy => Price::max(FIXED_PRECISION),
OrderSide::Sell => Price::min(FIXED_PRECISION),
};
let mut fills = if self.config.liquidity_consumption {
let size_prec = self.instrument.size_precision();
self.book
.get_all_crossed_levels(order.order_side(), price, size_prec)
} else {
let book_order = BookOrder::new(order.order_side(), price, order.quantity(), 0);
self.book.simulate_fills(&book_order)
};
if !self.fill_at_market
&& self.book_type == BookType::L1_MBP
&& !fills.is_empty()
&& matches!(
order.order_type(),
OrderType::StopMarket | OrderType::TrailingStopMarket | OrderType::MarketIfTouched
)
&& let Some(trigger_price) = order.trigger_price()
{
fills[0] = (trigger_price, fills[0].1);
let mut remaining_qty = order.leaves_qty();
let mut capped_fills = Vec::with_capacity(fills.len());
for (price, qty) in fills {
if remaining_qty.is_zero() {
break;
}
let mut capped_qty = qty.min(remaining_qty);
capped_qty.precision = qty.precision;
if capped_qty.is_zero() {
continue;
}
remaining_qty = remaining_qty - capped_qty;
capped_fills.push((price, capped_qty));
}
return capped_fills;
}
fills
}
fn determine_market_fill_model_price_and_volume(
&mut self,
order: &OrderAny,
) -> anyhow::Result<(Vec<(Price, Quantity)>, bool)> {
if let (Some(best_bid), Some(best_ask)) = (self.core.bid, self.core.ask)
&& let Some(book) = self.fill_model.get_orderbook_for_fill_simulation(
&self.instrument,
order,
best_bid,
best_ask,
)?
{
let price = match order.order_side() {
OrderSide::Buy => Price::max(FIXED_PRECISION),
OrderSide::Sell => Price::min(FIXED_PRECISION),
};
let book_order = BookOrder::new(order.order_side(), price, order.quantity(), 0);
let fills = book.simulate_fills(&book_order);
if !fills.is_empty() {
return Ok((fills, true));
}
}
Ok((self.determine_market_price_and_volume(order), false))
}
fn determine_limit_fill_model_price_and_volume(
&mut self,
order: &OrderAny,
) -> anyhow::Result<Vec<(Price, Quantity)>> {
if let (Some(best_bid), Some(best_ask)) = (self.core.bid, self.core.ask)
&& let Some(book) = self.fill_model.get_orderbook_for_fill_simulation(
&self.instrument,
order,
best_bid,
best_ask,
)?
&& let Some(limit_price) = order.price()
{
let book_order = BookOrder::new(order.order_side(), limit_price, order.quantity(), 0);
let fills = book.simulate_fills(&book_order);
if !fills.is_empty() {
return Ok(fills);
}
}
Ok(self.determine_limit_price_and_volume(order))
}
pub fn fill_market_order(&mut self, client_order_id: ClientOrderId) {
let mut order = match self.order_snapshot(client_order_id) {
Some(order) => order,
None => {
log::error!("Cannot fill market order: order {client_order_id} not found in cache");
return;
}
};
if order.is_closed() {
self.purge_stale_core_entry(client_order_id);
return;
}
if order.is_quote_quantity()
&& !self.instrument.is_inverse()
&& !self.convert_quote_to_base_quantity(&mut order)
{
return;
}
if let Some(filled_qty) = self.cached_filled_qty.get(&order.client_order_id())
&& filled_qty >= &order.quantity()
{
log::debug!(
"Ignoring fill as already filled pending application of events: {:?}, {:?}, {:?}, {:?}",
filled_qty,
order.quantity(),
order.filled_qty(),
order.quantity()
);
return;
}
let (venue_position_id, position) = self.fill_position_for_order(&order, Some(true));
if self.config.use_reduce_only && order.is_reduce_only() && position.is_none() {
log::warn!(
"Canceling REDUCE_ONLY {} as would increase position",
order.order_type()
);
self.cancel_order(&order, None);
return;
}
order.set_liquidity_side(LiquiditySide::Taker);
let (mut fills, from_synthetic) =
match self.determine_market_fill_model_price_and_volume(&order) {
Ok(result) => result,
Err(e) => {
log::error!(
"Cannot fill market order {}: fill model failed: {e}",
order.client_order_id()
);
return;
}
};
let protection_price: Option<Price> = if let Some(protection_points) =
self.config.price_protection_points
&& matches!(
order.order_type(),
OrderType::Market | OrderType::StopMarket
) {
protection_price_calculate(
self.instrument.price_increment(),
&order,
protection_points,
self.core.bid,
self.core.ask,
)
.ok()
} else {
None
};
if let Some(protection_price) = protection_price {
fills = self.filter_fills_by_protection(fills, &order, protection_price);
}
let is_trigger_price_fill = !self.fill_at_market
&& self.book_type == BookType::L1_MBP
&& matches!(
order.order_type(),
OrderType::StopMarket | OrderType::TrailingStopMarket | OrderType::MarketIfTouched
)
&& order.trigger_price().is_some();
if !from_synthetic && !is_trigger_price_fill {
fills = self.apply_liquidity_consumption(
fills,
order.order_side(),
order.leaves_qty(),
None,
);
}
if let Err(e) = self.apply_fills(
&order,
&fills,
LiquiditySide::Taker,
if self.config.use_reduce_only && order.is_reduce_only() {
venue_position_id
} else {
None
},
position.as_ref(),
protection_price,
) {
log::error!("Cannot fill market order {}: {e}", order.client_order_id());
}
}
fn filter_fills_by_protection(
&self,
fills: Vec<(Price, Quantity)>,
order: &OrderAny,
protection_price: Price,
) -> Vec<(Price, Quantity)> {
fills
.into_iter()
.filter(|(fill_price, _)| {
match order.order_side() {
OrderSide::Buy => *fill_price <= protection_price,
OrderSide::Sell => *fill_price >= protection_price,
}
})
.collect()
}
pub fn fill_limit_order(&mut self, client_order_id: ClientOrderId) {
let mut order = match self.order_snapshot(client_order_id) {
Some(order) => order,
None => {
log::error!("Cannot fill limit order: order {client_order_id} not found in cache");
return;
}
};
if order.is_closed() {
self.purge_stale_core_entry(client_order_id);
return;
}
if order.is_quote_quantity()
&& !self.instrument.is_inverse()
&& !self.convert_quote_to_base_quantity(&mut order)
{
return;
}
match order.price() {
Some(order_price) => {
let cached_filled_qty = self.cached_filled_qty.get(&order.client_order_id());
if let Some(&qty) = cached_filled_qty
&& qty >= order.quantity()
{
log::debug!(
"Ignoring fill as already filled pending application of events: {}, {}, {}, {}",
qty,
order.quantity(),
order.filled_qty(),
order.leaves_qty(),
);
return;
}
if order
.liquidity_side()
.is_some_and(|liquidity_side| liquidity_side == LiquiditySide::Maker)
{
let at_limit = if self.last_trade_size.is_some() && self.core.last.is_some() {
self.core.last.is_some_and(|last| last == order_price)
} else if order.order_side() == OrderSide::Buy {
self.core.bid.is_some_and(|bid| bid == order_price)
} else {
self.core.ask.is_some_and(|ask| ask == order_price)
};
if at_limit {
let is_limit_filled = match self.fill_model.is_limit_filled() {
Ok(value) => value,
Err(e) => {
log::error!(
"Cannot fill limit order {}: fill model failed: {e}",
order.client_order_id()
);
return;
}
};
if !is_limit_filled {
return; }
}
}
let queue_allowed_raw = if self.config.queue_position {
match self.determine_trade_fill_qty(&order) {
None | Some(0) => {
if matches!(order.time_in_force(), TimeInForce::Fok | TimeInForce::Ioc)
{
self.cancel_order(&order, None);
}
return;
}
Some(allowed) => Some(allowed),
}
} else {
None
};
let (venue_position_id, position) = self.fill_position_for_order(&order, None);
if self.config.use_reduce_only && order.is_reduce_only() && position.is_none() {
log::warn!(
"Canceling REDUCE_ONLY {} as would increase position",
order.order_type()
);
self.cancel_order(&order, None);
return;
}
let tc_before = self.trade_consumption;
let mut fills = match self.determine_limit_fill_model_price_and_volume(&order) {
Ok(fills) => fills,
Err(e) => {
log::error!(
"Cannot fill limit order {}: fill model failed: {e}",
order.client_order_id()
);
return;
}
};
if let Some(allowed_raw) = queue_allowed_raw {
let size_prec = self.instrument.size_precision();
let mut remaining = allowed_raw;
fills = fills
.into_iter()
.filter_map(|(price, qty)| {
if remaining == 0 {
return None;
}
let capped = qty.raw().min(remaining);
remaining -= capped;
Some((price, Quantity::from_raw(capped, size_prec)))
})
.collect();
let consumed: QuantityRaw = fills.iter().map(|(_, qty)| qty.raw()).sum();
if let Some(excess) = self.queue_excess.get_mut(&order.client_order_id()) {
*excess = excess.saturating_sub(consumed);
}
self.trade_consumption = tc_before + consumed;
}
if fills.is_empty() && self.config.liquidity_consumption {
log::debug!(
"Skipping fill for {}: no liquidity available after consumption",
order.client_order_id()
);
if matches!(order.time_in_force(), TimeInForce::Fok | TimeInForce::Ioc) {
self.cancel_order(&order, None);
}
return;
}
let liquidity_side = order.liquidity_side().unwrap();
if let Err(e) = self.apply_fills(
&order,
&fills,
liquidity_side,
venue_position_id,
position.as_ref(),
None,
) {
log::error!("Cannot fill limit order {}: {e}", order.client_order_id());
}
}
None => panic!("Limit order must have a price"),
}
}
fn fill_position_for_order(
&mut self,
order: &OrderAny,
generate: Option<bool>,
) -> (Option<PositionId>, Option<Position>) {
if self.oms_type == OmsType::Hedging
&& self.config.use_reduce_only
&& order.is_reduce_only()
{
let cache = self.cache.as_ref().borrow();
if let Some(position) = cache.position_for_order(&order.client_order_id()) {
let position = position.clone_without_events();
return (Some(position.id), Some(position));
}
if let Some(position) = Self::open_position_reduced_by_order(&cache, order) {
return (Some(position.id), Some(position));
}
}
let venue_position_id = self.ids_generator.get_position_id(order, generate);
let position = {
let cache = self.cache.as_ref().borrow();
venue_position_id
.as_ref()
.and_then(|position_id| cache.position(position_id))
.map(|position| position.clone_without_events())
};
(venue_position_id, position)
}
fn position_for_order_in_cache(&self, cache: &Cache, order: &OrderAny) -> Option<Position> {
if let Some(position) = cache.position_for_order(&order.client_order_id()) {
return Some(position.clone_without_events());
}
if self.oms_type == OmsType::Netting {
let position_id = PositionId::new(
format!("{}-{}", order.instrument_id(), order.strategy_id()).as_str(),
);
return cache
.position(&position_id)
.map(|position| position.clone_without_events());
}
if self.oms_type == OmsType::Hedging
&& self.config.use_reduce_only
&& order.is_reduce_only()
{
return Self::open_position_reduced_by_order(cache, order);
}
None
}
fn open_position_reduced_by_order(cache: &Cache, order: &OrderAny) -> Option<Position> {
cache
.positions_open(
None,
Some(&order.instrument_id()),
Some(&order.strategy_id()),
None,
None,
)
.into_iter()
.find(|position| order.would_reduce_only(position.side, position.quantity))
.map(|position| position.clone_without_events())
}
fn apply_fills(
&mut self,
order: &OrderAny,
fills: &[(Price, Quantity)],
liquidity_side: LiquiditySide,
venue_position_id: Option<PositionId>,
position: Option<&Position>,
protection_price: Option<Price>,
) -> anyhow::Result<()> {
if order.time_in_force() == TimeInForce::Fok {
let mut total_size = Quantity::zero(order.quantity().precision);
for &(fill_px, fill_qty) in fills {
if self
.normalize_price_for_current_instrument(fill_px)
.is_some()
&& let Some(fill_qty) = self.normalize_quantity_for_current_instrument(fill_qty)
{
total_size = total_size.add(fill_qty);
}
}
if order.leaves_qty() > total_size {
self.cancel_order(order, None);
return Ok(());
}
}
if fills.is_empty() {
if order.status() == OrderStatus::Submitted {
self.generate_order_rejected(
order,
format!("No market for {}", order.instrument_id()).into(),
);
} else {
log::error!(
"Cannot fill order: no fills from book when fills were expected (check size in data)"
);
return Ok(());
}
}
let venue_position_id = if self.oms_type == OmsType::Netting {
None
} else {
venue_position_id
};
let mut initial_market_to_limit_fill = false;
let mut total_filled = self
.cached_filled_qty
.get(&order.client_order_id())
.copied()
.unwrap_or_else(|| order.filled_qty());
let initial_total_filled = total_filled;
let mut last_fill_px: Option<Price> = None;
let mut reduce_only_remaining = None;
let mut reduce_only_filled = None;
if self.config.use_reduce_only
&& order.is_reduce_only()
&& let Some(current_position) = position
{
let remaining = self.position_quantity_remaining(order, current_position)?;
if remaining.is_zero() {
self.cancel_order(order, None);
return Ok(());
}
reduce_only_remaining = Some(remaining);
reduce_only_filled = Some(total_filled);
}
for &(fill_px, fill_qty) in fills {
let Some(mut fill_px) = self.normalize_fill_price(fill_px, order.client_order_id())
else {
continue;
};
let Some(fill_qty) = self.normalize_fill_quantity(fill_qty, order.client_order_id())
else {
continue;
};
if order.filled_qty() == Quantity::zero(order.filled_qty().precision)
&& order.order_type() == OrderType::MarketToLimit
{
self.generate_order_updated(order, order.quantity(), Some(fill_px), None, None);
initial_market_to_limit_fill = true;
}
if self.book_type == BookType::L1_MBP && self.fill_model.is_slipped()? {
fill_px = match order.order_side() {
OrderSide::Buy => fill_px.add(self.instrument.price_increment()),
OrderSide::Sell => fill_px.sub(self.instrument.price_increment()),
}
}
let mut effective_fill_qty = fill_qty;
if let Some(remaining) = reduce_only_remaining {
if remaining.is_zero() {
return Ok(());
}
if effective_fill_qty > remaining {
let precision = effective_fill_qty.precision;
effective_fill_qty = remaining;
effective_fill_qty.precision = precision;
}
}
if fill_qty.is_zero() {
if fills.len() == 1 && order.status() == OrderStatus::Submitted {
self.generate_order_rejected(
order,
format!("No market for {}", order.instrument_id()).into(),
);
}
return Ok(());
}
let capped_fill_qty = min(
effective_fill_qty,
order.quantity().saturating_sub(total_filled),
);
let reduce_only_exhausts_position =
reduce_only_remaining.is_some_and(|remaining| capped_fill_qty >= remaining);
if reduce_only_exhausts_position {
let mut reduce_only_target = reduce_only_filled
.unwrap_or(initial_total_filled)
.checked_add(capped_fill_qty)
.expect("Overflow occurred when adding reduce-only target quantity");
reduce_only_target.precision = order.quantity().precision;
if order.quantity() != reduce_only_target {
self.generate_order_updated(order, reduce_only_target, None, None, None);
}
}
total_filled = total_filled.add(capped_fill_qty);
if let Some(remaining) = reduce_only_remaining.as_mut() {
*remaining = *remaining - capped_fill_qty.min(*remaining);
}
if let Some(filled) = reduce_only_filled.as_mut() {
*filled = filled
.checked_add(capped_fill_qty)
.expect("Overflow occurred when adding reduce-only filled quantity");
}
self.fill_order(
order,
fill_px,
effective_fill_qty,
liquidity_side,
venue_position_id,
position,
)?;
last_fill_px = Some(fill_px);
if order.order_type() == OrderType::MarketToLimit && initial_market_to_limit_fill {
return Ok(());
}
if reduce_only_exhausts_position {
self.purge_cached_filled_qty_if_closed(order.client_order_id());
return Ok(());
}
}
let leaves_remaining = total_filled < order.quantity();
let filled_in_loop = total_filled > initial_total_filled;
if order.time_in_force() == TimeInForce::Ioc && leaves_remaining {
self.cancel_order(order, None);
return Ok(());
}
if leaves_remaining
&& (order.is_open() || filled_in_loop)
&& self.book_type == BookType::L1_MBP
&& matches!(
order.order_type(),
OrderType::Market
| OrderType::MarketIfTouched
| OrderType::StopMarket
| OrderType::TrailingStopMarket
)
{
let Some(last_fill_px) = last_fill_px else {
return Ok(());
};
let side = order.order_side();
let slip_fill_px = match side {
OrderSide::Buy => last_fill_px.add(self.instrument.price_increment()),
OrderSide::Sell => last_fill_px.sub(self.instrument.price_increment()),
};
if let Some(protection_price) = protection_price {
let exceeds_boundary = match side {
OrderSide::Buy => slip_fill_px > protection_price,
OrderSide::Sell => slip_fill_px < protection_price,
};
if exceeds_boundary {
return Ok(());
}
}
let mut leaves_qty = order.quantity().saturating_sub(total_filled);
if let Some(remaining) = reduce_only_remaining {
if remaining.is_zero() {
return Ok(());
}
if leaves_qty > remaining {
let precision = leaves_qty.precision;
leaves_qty = remaining;
leaves_qty.precision = precision;
}
if leaves_qty >= remaining {
let mut reduce_only_target = reduce_only_filled
.unwrap_or(initial_total_filled)
.checked_add(leaves_qty)
.expect("Overflow occurred when adding reduce-only target quantity");
reduce_only_target.precision = order.quantity().precision;
if order.quantity() != reduce_only_target {
self.generate_order_updated(order, reduce_only_target, None, None, None);
}
}
}
if leaves_qty.is_zero() {
return Ok(());
}
self.fill_order(
order,
slip_fill_px,
leaves_qty,
liquidity_side,
venue_position_id,
position,
)?;
self.purge_cached_filled_qty_if_closed(order.client_order_id());
}
Ok(())
}
fn normalize_fill_price(
&self,
fill_px: Price,
client_order_id: ClientOrderId,
) -> Option<Price> {
let normalized = self.normalize_price_for_current_instrument(fill_px);
if normalized.is_none() {
log::warn!(
"Skipping fill for {client_order_id}: fill price {fill_px} is not compatible \
with {} price_precision={} price_increment={}",
self.instrument.id(),
self.instrument.price_precision(),
self.instrument.price_increment()
);
}
normalized
}
fn normalize_fill_quantity(
&self,
fill_qty: Quantity,
client_order_id: ClientOrderId,
) -> Option<Quantity> {
let normalized = self.normalize_quantity_for_current_instrument(fill_qty);
if normalized.is_none() {
log::warn!(
"Skipping fill for {client_order_id}: fill quantity {fill_qty} is not compatible \
with {} size_precision={}",
self.instrument.id(),
self.instrument.size_precision()
);
}
normalized
}
fn position_quantity_remaining(
&mut self,
order: &OrderAny,
position: &Position,
) -> anyhow::Result<Quantity> {
self.purge_applied_fills();
let mut quantity = match position.side {
PositionSide::Long => position.quantity.as_decimal(),
PositionSide::Short => -position.quantity.as_decimal(),
PositionSide::Flat => Decimal::ZERO,
};
for fill in self.pending_fills.values() {
if fill.position_id == Some(position.id) {
quantity = quantity
.checked_add(fill.quantity_change)
.ok_or_else(|| anyhow::anyhow!("Pending position quantity overflow"))?;
}
}
if (order.is_buy() && quantity >= Decimal::ZERO)
|| (order.is_sell() && quantity <= Decimal::ZERO)
{
return Ok(Quantity::zero(position.quantity.precision));
}
Ok(Quantity::from_decimal_dp(
quantity.abs(),
position.quantity.precision,
)?)
}
fn purge_applied_fills(&mut self) {
let cache = self.cache.borrow();
self.pending_fills.retain(|trade_id, fill| {
fill.position_id = fill
.position_id
.or_else(|| cache.position_id(&fill.client_order_id).copied());
let Some(position_id) = fill.position_id else {
return cache.order_exists(&fill.client_order_id);
};
let Some(position) = cache.position(&position_id) else {
return cache.order_exists(&fill.client_order_id);
};
if position.trade_ids.contains(trade_id) {
return false;
}
let opening_trade_id = position.events.first().map(|event| event.trade_id);
if opening_trade_id != fill.opening_trade_id {
if position.replay_events.iter().any(|event| {
matches!(event, PositionReplayEvent::Filled(event) if event.trade_id == *trade_id)
}) || cache.position_snapshots(Some(&position_id), None).iter()
.any(|snapshot| snapshot.trade_ids.contains(trade_id))
{
return false;
}
fill.opening_trade_id = opening_trade_id;
}
true
});
}
fn fill_order(
&mut self,
order: &OrderAny,
last_px: Price,
last_qty: Quantity,
liquidity_side: LiquiditySide,
venue_position_id: Option<PositionId>,
position: Option<&Position>,
) -> anyhow::Result<()> {
self.check_size_precision(last_qty.precision, "fill quantity")?;
let (last_qty, new_filled_qty) =
if let Some(filled_qty) = self.cached_filled_qty.get(&order.client_order_id()) {
let leaves_qty = order.quantity().saturating_sub(*filled_qty);
let last_qty = min(last_qty, leaves_qty);
(last_qty, *filled_qty + last_qty)
} else {
let last_qty = min(last_qty, order.quantity());
(last_qty, last_qty)
};
if last_qty.is_zero() {
return Ok(());
}
let fee_order;
let commission_order = {
let mut cloned = order.clone();
write_filled_qty(&mut cloned, new_filled_qty.saturating_sub(last_qty));
if order.liquidity_side() != Some(liquidity_side) {
cloned.set_liquidity_side(liquidity_side);
}
fee_order = cloned;
&fee_order
};
let underlying_px = self.fee_underlying_price()?;
let commission = self.fee_model.get_commission_with_context(
commission_order,
last_qty,
last_px,
&self.instrument,
underlying_px,
)?;
let reduce_only_order_ids = position
.map(|position| self.reduce_only_order_ids(position.id))
.unwrap_or_default();
self.cached_filled_qty
.insert(order.client_order_id(), new_filled_qty);
let venue_order_id = self.ids_generator.get_venue_order_id(order).unwrap();
self.generate_order_filled(
order,
venue_order_id,
venue_position_id,
last_qty,
last_px,
self.instrument.quote_currency(),
commission,
liquidity_side,
);
let post_fill_filled_qty = self
.cached_filled_qty
.get(&order.client_order_id())
.copied()
.unwrap_or(order.filled_qty());
let post_fill_leaves_qty = order.quantity().saturating_sub(post_fill_filled_qty);
let fully_filled = post_fill_leaves_qty.is_zero();
if order.is_closed() || fully_filled {
if self.core.order_exists(order.client_order_id()) {
self.delete_core_order(order.client_order_id());
}
self.remove_queue_position(order.client_order_id());
if order.order_type() != OrderType::MarketToLimit {
self.purge_cached_filled_qty_if_closed(order.client_order_id());
}
}
if self.config.support_contingent_orders
&& let Some(contingency_type) = order.contingency_type()
{
match contingency_type {
ContingencyType::Oto => {
if let Some(linked_orders_ids) = order.linked_order_ids() {
for client_order_id in linked_orders_ids {
let mut child_order = match self.order_snapshot(*client_order_id) {
Some(child_order) => child_order,
None => anyhow::bail!("Order {client_order_id} not found in cache"),
};
if child_order.is_closed() || child_order.is_active_local() {
continue;
}
if self.inflight_orders.contains(*client_order_id) {
continue;
}
if let (None, Some(position_id)) =
(child_order.position_id(), order.position_id())
{
self.cache
.borrow_mut()
.add_position_id(
&position_id,
&self.venue,
client_order_id,
&child_order.strategy_id(),
)
.unwrap();
log::debug!(
"Added position id {position_id} to cache for order {client_order_id}"
);
}
if (!child_order.is_open())
|| (matches!(child_order.status(), OrderStatus::PendingUpdate)
&& child_order
.previous_status()
.is_some_and(|s| matches!(s, OrderStatus::Submitted)))
{
let account_id = order
.account_id()
.or_else(|| self.account_ids.get(&order.trader_id()).copied())
.ok_or_else(|| {
anyhow::anyhow!(
"Account ID not found for trader {}",
order.trader_id()
)
})?;
self.process_order(&mut child_order, account_id);
}
}
} else {
log::error!(
"OTO order {} does not have linked orders",
order.client_order_id()
);
}
}
ContingencyType::Oco => {
if let Some(linked_orders_ids) = order.linked_order_ids() {
for client_order_id in linked_orders_ids {
let child_order = match self.order_snapshot(*client_order_id) {
Some(child_order) => child_order,
None => anyhow::bail!("Order {client_order_id} not found in cache"),
};
if child_order.is_closed() || child_order.is_active_local() {
continue;
}
self.cancel_order(&child_order, Some(false));
}
} else {
log::error!(
"OCO order {} does not have linked orders",
order.client_order_id()
);
}
}
ContingencyType::Ouo => {
if let Some(linked_orders_ids) = order.linked_order_ids() {
for client_order_id in linked_orders_ids {
let child_order = match self.order_snapshot(*client_order_id) {
Some(child_order) => child_order,
None => anyhow::bail!("Order {client_order_id} not found in cache"),
};
if child_order.is_active_local() {
continue;
}
let child_filled_qty = self
.cached_filled_qty
.get(&child_order.client_order_id())
.copied()
.unwrap_or(child_order.filled_qty());
if post_fill_leaves_qty.is_zero() && child_order.is_open() {
self.cancel_order(&child_order, None);
} else if child_order.is_open()
&& child_filled_qty >= post_fill_leaves_qty
{
self.cancel_order(&child_order, Some(false));
} else if post_fill_leaves_qty.non_zero()
&& post_fill_leaves_qty != child_order.leaves_qty()
{
let price = child_order.price();
let trigger_price = child_order.trigger_price();
self.update_order(
&child_order,
Some(post_fill_leaves_qty),
price,
trigger_price,
Some(false),
);
}
}
} else {
log::error!(
"OUO order {} does not have linked orders",
order.client_order_id()
);
}
}
}
}
if let Some(position) = position {
let mut reduce_only_order_ids = reduce_only_order_ids;
reduce_only_order_ids.extend(self.reduce_only_order_ids(position.id));
reduce_only_order_ids.sort_unstable();
reduce_only_order_ids.dedup();
self.sync_reduce_only_orders(order, position, &reduce_only_order_ids)?;
}
Ok(())
}
fn reduce_only_order_ids(&self, position_id: PositionId) -> Vec<ClientOrderId> {
if !self.config.use_reduce_only {
return Vec::new();
}
let cache = self.cache.borrow();
let mut order_ids = Vec::new();
for resting in self.core.iter_orders() {
let Some(order) = cache.order(&resting.client_order_id) else {
continue;
};
if !order.is_reduce_only() || !order.is_open() || !order.is_passive() {
continue;
}
let matches_position = match cache.position_id(&resting.client_order_id) {
Some(id) => *id == position_id,
None => self
.position_for_order_in_cache(&cache, &order)
.is_some_and(|position| position.id == position_id),
};
if matches_position {
order_ids.push(resting.client_order_id);
}
}
order_ids.sort_unstable();
order_ids
}
fn sync_reduce_only_orders(
&mut self,
filled_order: &OrderAny,
position: &Position,
order_ids: &[ClientOrderId],
) -> anyhow::Result<()> {
for &client_order_id in order_ids {
if client_order_id == filled_order.client_order_id()
|| !self.core.order_exists(client_order_id)
{
continue;
}
let Some(order) = self.order_snapshot(client_order_id) else {
continue;
};
if !order.is_reduce_only() || !order.is_open() || !order.is_passive() {
continue;
}
let position = self.cache.borrow().position(&position.id).map_or_else(
|| position.clone_without_events(),
|position| position.clone_without_events(),
);
let remaining = self.position_quantity_remaining(&order, &position)?;
if remaining.is_zero() {
self.cancel_reduce_only_order(&order, filled_order.client_order_id())?;
continue;
}
let leaves = self.parent_capped_leaves(&order, remaining);
let target = order.filled_qty().checked_add(leaves).ok_or_else(|| {
anyhow::anyhow!("Reduce-only quantity overflow for order {client_order_id}")
})?;
if order.quantity() != target {
self.generate_order_updated(
&order,
target,
order.price(),
order.trigger_price(),
None,
);
if target == order.filled_qty() {
self.cancel_reduce_only_order(&order, filled_order.client_order_id())?;
} else if self.config.support_contingent_orders
&& order.contingency_type() == Some(ContingencyType::Ouo)
{
self.sync_ouo_leaves(&order, leaves, filled_order.client_order_id())?;
}
}
}
Ok(())
}
fn cancel_reduce_only_order(
&mut self,
order: &OrderAny,
filled_order_id: ClientOrderId,
) -> anyhow::Result<()> {
let propagate = self.config.support_contingent_orders
&& order.contingency_type() == Some(ContingencyType::Ouo);
self.cancel_order(order, Some(!propagate));
if propagate {
self.sync_ouo_leaves(
order,
Quantity::zero(order.quantity().precision),
filled_order_id,
)?;
}
Ok(())
}
fn parent_capped_leaves(&self, order: &OrderAny, leaves: Quantity) -> Quantity {
let parent = if self.config.support_contingent_orders {
order
.parent_order_id()
.and_then(|id| self.order_snapshot(id))
} else {
None
};
parent.map_or(leaves, |parent| {
min(
leaves,
parent.filled_qty().saturating_sub(order.filled_qty()),
)
})
}
fn sync_ouo_leaves(
&mut self,
order: &OrderAny,
leaves: Quantity,
filled_order_id: ClientOrderId,
) -> anyhow::Result<()> {
for &client_order_id in order.linked_order_ids().into_iter().flatten() {
if client_order_id == filled_order_id || !self.core.order_exists(client_order_id) {
continue;
}
let Some(sibling) = self.order_snapshot(client_order_id) else {
continue;
};
if sibling.is_closed() || sibling.is_active_local() || !sibling.is_passive() {
continue;
}
if leaves.is_zero() {
self.cancel_order(&sibling, Some(false));
continue;
}
if !sibling.is_open() {
continue;
}
let leaves = self.parent_capped_leaves(&sibling, leaves);
let target = sibling.filled_qty().checked_add(leaves).ok_or_else(|| {
anyhow::anyhow!("OUO quantity overflow for order {client_order_id}")
})?;
if sibling.quantity() != target {
self.generate_order_updated(
&sibling,
target,
sibling.price(),
sibling.trigger_price(),
None,
);
}
if leaves.is_zero() {
self.cancel_order(&sibling, Some(false));
}
}
Ok(())
}
fn fee_underlying_price(&self) -> CorrectnessResult<Option<Price>> {
if !matches!(
self.instrument,
InstrumentAny::CryptoOption(_) | InstrumentAny::OptionContract(_)
) {
return Ok(None);
}
let Some(underlying) = self.instrument.underlying() else {
return Ok(None);
};
let underlying_id = InstrumentId::from(format!("{underlying}.{}", self.venue).as_str());
let instrument_id = self.instrument.id();
let cache = self.cache.borrow();
if let Some(price) = cache
.price(&underlying_id, PriceType::Last)
.or_else(|| cache.price(&underlying_id, PriceType::Mark))
.or_else(|| cache.price(&underlying_id, PriceType::Mid))
{
return Ok(Some(price));
}
cache
.option_greeks(&instrument_id)
.and_then(|greeks| greeks.underlying_price)
.map(|price| Price::new_checked(price, FIXED_PRECISION))
.transpose()
}
fn cached_order_is_closed(&self, client_order_id: ClientOrderId) -> bool {
self.cache
.borrow()
.order(&client_order_id)
.is_none_or(|order| order.is_closed())
}
fn purge_cached_filled_qty_if_closed(&mut self, client_order_id: ClientOrderId) {
if self.cached_order_is_closed(client_order_id) {
self.cached_filled_qty.swap_remove(&client_order_id);
}
}
fn purge_closed_cached_filled_qty(&mut self) {
let client_order_ids: Vec<ClientOrderId> = self.cached_filled_qty.keys().copied().collect();
for client_order_id in client_order_ids {
self.purge_cached_filled_qty_if_closed(client_order_id);
}
}
fn update_limit_order(
&mut self,
order: &OrderAny,
quantity: Quantity,
price: Price,
) -> ModifyOutcome {
if self.core.is_limit_matched(order.order_side(), price) {
if order.is_post_only() {
self.generate_order_modify_rejected(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
Ustr::from(format!(
"POST_ONLY {} {} order with new limit px of {} would have been a TAKER: bid={}, ask={}",
order.order_type(),
order.order_side(),
price,
self.core.bid.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core.ask.map_or_else(|| "None".to_string(), |p| p.to_string())
).as_str()),
order.venue_order_id(),
order.account_id(),
);
return ModifyOutcome::Rejected;
}
self.generate_order_updated(order, quantity, Some(price), None, None);
let client_order_id = order.client_order_id();
if let Some(mut order) = self.cache.borrow_mut().order_mut(&client_order_id) {
order.set_liquidity_side(LiquiditySide::Taker);
}
self.fill_limit_order(client_order_id);
return ModifyOutcome::Applied;
}
self.generate_order_updated(order, quantity, Some(price), None, None);
ModifyOutcome::Applied
}
fn update_stop_market_order(
&self,
order: &OrderAny,
quantity: Quantity,
trigger_price: Price,
) -> ModifyOutcome {
if self.core.is_stop_matched_with_trigger_type(
order.order_side(),
trigger_price,
order.trigger_type().unwrap_or(TriggerType::Default),
) {
self.generate_order_modify_rejected(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
Ustr::from(
format!(
"{} {} order new stop px of {} was in the market: bid={}, ask={}",
order.order_type(),
order.order_side(),
trigger_price,
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
)
.as_str(),
),
order.venue_order_id(),
order.account_id(),
);
return ModifyOutcome::Rejected;
}
self.generate_order_updated(order, quantity, None, Some(trigger_price), None);
ModifyOutcome::Applied
}
fn update_stop_limit_order(
&mut self,
order: &OrderAny,
quantity: Quantity,
price: Price,
trigger_price: Price,
) -> ModifyOutcome {
if order.is_triggered().is_some_and(|t| t) {
if self.core.is_limit_matched(order.order_side(), price) {
return self.update_limit_order(order, quantity, price);
}
} else {
if self.core.is_stop_matched_with_trigger_type(
order.order_side(),
trigger_price,
order.trigger_type().unwrap_or(TriggerType::Default),
) {
self.generate_order_modify_rejected(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
Ustr::from(
format!(
"{} {} order new stop px of {} was in the market: bid={}, ask={}",
order.order_type(),
order.order_side(),
trigger_price,
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
)
.as_str(),
),
order.venue_order_id(),
order.account_id(),
);
return ModifyOutcome::Rejected;
}
}
self.generate_order_updated(order, quantity, Some(price), Some(trigger_price), None);
ModifyOutcome::Applied
}
fn update_market_if_touched_order(
&self,
order: &OrderAny,
quantity: Quantity,
trigger_price: Price,
) -> ModifyOutcome {
if self.core.is_touch_triggered_with_trigger_type(
order.order_side(),
trigger_price,
order.trigger_type().unwrap_or(TriggerType::Default),
) {
self.generate_order_modify_rejected(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
Ustr::from(
format!(
"{} {} order new trigger px of {} was in the market: bid={}, ask={}",
order.order_type(),
order.order_side(),
trigger_price,
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
)
.as_str(),
),
order.venue_order_id(),
order.account_id(),
);
return ModifyOutcome::Rejected;
}
self.generate_order_updated(order, quantity, None, Some(trigger_price), None);
ModifyOutcome::Applied
}
fn update_limit_if_touched_order(
&mut self,
order: &OrderAny,
quantity: Quantity,
price: Price,
trigger_price: Price,
) -> ModifyOutcome {
if order.is_triggered().is_some_and(|t| t) {
if self.core.is_limit_matched(order.order_side(), price) {
return self.update_limit_order(order, quantity, price);
}
} else {
if self.core.is_touch_triggered_with_trigger_type(
order.order_side(),
trigger_price,
order.trigger_type().unwrap_or(TriggerType::Default),
) {
self.generate_order_modify_rejected(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
Ustr::from(
format!(
"{} {} order new trigger px of {} was in the market: bid={}, ask={}",
order.order_type(),
order.order_side(),
trigger_price,
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
)
.as_str(),
),
order.venue_order_id(),
order.account_id(),
);
return ModifyOutcome::Rejected;
}
}
self.generate_order_updated(order, quantity, Some(price), Some(trigger_price), None);
ModifyOutcome::Applied
}
fn update_trailing_stop_order(&self, order: &OrderAny) {
let (new_trigger_price, new_price) = match trailing_stop_calculate(
self.instrument.price_increment(),
order.trigger_price(),
order,
self.core.bid,
self.core.ask,
self.core.last,
) {
Ok(prices) => prices,
Err(e) => {
log::debug!("Cannot calculate trailing-stop update: {e}");
return;
}
};
if new_trigger_price.is_none() && new_price.is_none() {
return;
}
self.generate_order_updated(order, order.quantity(), new_price, new_trigger_price, None);
}
fn accept_order(&mut self, order: &mut OrderAny) {
if order.is_closed() {
return;
}
if order.status() != OrderStatus::Accepted {
let venue_order_id = self.ids_generator.get_venue_order_id(order).unwrap();
let event = self.create_order_accepted(order, venue_order_id);
if let Err(e) = order.apply(event.clone()) {
log::warn!(
"Skipping local apply of accepted event for {}: {e}",
order.client_order_id(),
);
}
self.dispatch_order_event(event);
if matches!(
order.order_type(),
OrderType::TrailingStopLimit | OrderType::TrailingStopMarket
) && order.trigger_price().is_none()
&& self.maybe_activate_trailing_stop(
order,
self.core.bid,
self.core.ask,
self.core.last,
)
{
self.update_trailing_stop_order(order);
}
}
let match_info = Self::matching_core_entry(order);
self.track_post_match_order(order);
self.core.add_order(match_info);
}
fn track_post_match_order(&mut self, order: &OrderAny) {
self.post_match_order_ids.insert(order.client_order_id());
}
fn delete_core_order(&mut self, client_order_id: ClientOrderId) {
self.post_match_order_ids.swap_remove(&client_order_id);
let _ = self.core.delete_order(client_order_id);
}
fn requires_post_match_maintenance(order: &OrderAny) -> bool {
order.expire_time().is_some()
|| matches!(
order.order_type(),
OrderType::TrailingStopMarket | OrderType::TrailingStopLimit
)
}
fn matching_core_entry(order: &OrderAny) -> RestingOrder {
let triggered_limit_style = matches!(
order.order_type(),
OrderType::StopLimit | OrderType::LimitIfTouched | OrderType::TrailingStopLimit
) && order.is_triggered().is_some_and(|triggered| triggered);
RestingOrder::new_with_trigger_type(
order.client_order_id(),
order.order_side(),
order.order_type(),
Some(order.trigger_type().unwrap_or(TriggerType::Default)),
if triggered_limit_style {
None
} else {
order.trigger_price()
},
order.price(),
match order {
OrderAny::TrailingStopMarket(o) => o.is_activated,
OrderAny::TrailingStopLimit(o) => o.is_activated,
_ => true,
},
)
}
fn expire_order(&mut self, order: &OrderAny) {
self.remove_queue_position(order.client_order_id());
if self.config.support_contingent_orders && order.contingency_type().is_some() {
self.cancel_contingent_orders(order, &[]);
}
self.generate_order_expired(order);
}
fn cancel_order(&mut self, order: &OrderAny, cancel_contingencies: Option<bool>) {
self.cancel_order_excluding(order, cancel_contingencies, &[]);
}
fn cancel_order_excluding(
&mut self,
order: &OrderAny,
cancel_contingencies: Option<bool>,
excluded: &[ClientOrderId],
) {
if self.inflight_orders.contains(order.client_order_id()) {
return;
}
let cancel_contingencies = cancel_contingencies.unwrap_or(true);
if order.is_active_local()
&& !matches!(
(order.status(), order.order_type(), order.time_in_force()),
(
OrderStatus::Initialized | OrderStatus::Released,
OrderType::Market,
TimeInForce::Ioc | TimeInForce::Fok
)
)
{
log::error!(
"Cannot cancel an order with {} from the matching engine",
order.status()
);
return;
}
if self.core.order_exists(order.client_order_id()) {
self.delete_core_order(order.client_order_id());
}
self.remove_queue_position(order.client_order_id());
self.cached_filled_qty.swap_remove(&order.client_order_id());
let venue_order_id = self.ids_generator.get_venue_order_id(order).unwrap();
self.generate_order_canceled(order, venue_order_id);
if self.config.support_contingent_orders
&& order.contingency_type().is_some()
&& cancel_contingencies
{
self.cancel_contingent_orders(order, excluded);
}
}
fn update_order(
&mut self,
order: &OrderAny,
quantity: Option<Quantity>,
price: Option<Price>,
trigger_price: Option<Price>,
update_contingencies: Option<bool>,
) -> bool {
if self.inflight_orders.contains(order.client_order_id()) {
return false;
}
let update_contingencies = update_contingencies.unwrap_or(true);
let quantity = quantity.unwrap_or(order.quantity());
let price_prec = self.instrument.price_precision();
let size_prec = self.instrument.size_precision();
let instrument_id = self.instrument.id();
if quantity.precision != size_prec {
self.generate_order_modify_rejected(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
Ustr::from(&format!(
"Invalid update quantity precision {}, expected {size_prec} for {instrument_id}",
quantity.precision
)),
order.venue_order_id(),
order.account_id(),
);
return false;
}
if let Some(px) = price
&& px.precision != price_prec
{
self.generate_order_modify_rejected(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
Ustr::from(&format!(
"Invalid update price precision {}, expected {price_prec} for {instrument_id}",
px.precision
)),
order.venue_order_id(),
order.account_id(),
);
return false;
}
if let Some(tp) = trigger_price
&& tp.precision != price_prec
{
self.generate_order_modify_rejected(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
Ustr::from(&format!(
"Invalid update trigger_price precision {}, expected {price_prec} for {instrument_id}",
tp.precision
)),
order.venue_order_id(),
order.account_id(),
);
return false;
}
let filled_qty = self
.cached_filled_qty
.get(&order.client_order_id())
.copied()
.unwrap_or(order.filled_qty());
if quantity < filled_qty {
self.generate_order_modify_rejected(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
Ustr::from(&format!(
"Cannot reduce order quantity {quantity} below filled quantity {filled_qty}",
)),
order.venue_order_id(),
order.account_id(),
);
return false;
}
let outcome = match order {
OrderAny::Limit(_) | OrderAny::MarketToLimit(_) => {
let price = price.unwrap_or(order.price().unwrap());
self.update_limit_order(order, quantity, price)
}
OrderAny::StopMarket(_) => {
let trigger_price = trigger_price.unwrap_or(order.trigger_price().unwrap());
self.update_stop_market_order(order, quantity, trigger_price)
}
OrderAny::StopLimit(_) => {
let price = price.unwrap_or(order.price().unwrap());
let trigger_price = trigger_price.unwrap_or(order.trigger_price().unwrap());
self.update_stop_limit_order(order, quantity, price, trigger_price)
}
OrderAny::MarketIfTouched(_) => {
let trigger_price = trigger_price.unwrap_or(order.trigger_price().unwrap());
self.update_market_if_touched_order(order, quantity, trigger_price)
}
OrderAny::LimitIfTouched(_) => {
let price = price.unwrap_or(order.price().unwrap());
let trigger_price = trigger_price.unwrap_or(order.trigger_price().unwrap());
self.update_limit_if_touched_order(order, quantity, price, trigger_price)
}
OrderAny::TrailingStopMarket(_) => {
if let Some(trigger_price) = trigger_price.or(order.trigger_price()) {
self.update_market_if_touched_order(order, quantity, trigger_price)
} else {
self.generate_order_updated(order, quantity, None, trigger_price, None);
ModifyOutcome::Applied
}
}
OrderAny::TrailingStopLimit(_) => {
match (
price.or(order.price()),
trigger_price.or(order.trigger_price()),
) {
(Some(price), Some(trigger_price)) => {
self.update_limit_if_touched_order(order, quantity, price, trigger_price)
}
_ => {
self.generate_order_updated(order, quantity, price, trigger_price, None);
ModifyOutcome::Applied
}
}
}
_ => {
panic!(
"Unsupported order type {} for update_order",
order.order_type()
);
}
};
if outcome == ModifyOutcome::Rejected {
return false;
}
let new_leaves_qty = quantity.saturating_sub(filled_qty);
if new_leaves_qty.is_zero() {
if self.config.support_contingent_orders
&& order.contingency_type().is_some()
&& update_contingencies
{
self.update_contingent_order(order, quantity);
}
self.cancel_order(order, Some(false));
return true;
}
if self.config.support_contingent_orders
&& order.contingency_type().is_some()
&& update_contingencies
{
self.update_contingent_order(order, quantity);
}
true
}
pub fn trigger_stop_order(&mut self, client_order_id: ClientOrderId) {
let order = match self.order_snapshot(client_order_id) {
Some(order) => order,
None => {
log::error!(
"Cannot trigger stop order: order {client_order_id} not found in cache"
);
return;
}
};
if order.is_closed() {
log::debug!("Cannot trigger stop order: {client_order_id} already closed");
return;
}
match order.order_type() {
OrderType::StopLimit | OrderType::LimitIfTouched | OrderType::TrailingStopLimit => {
self.trigger_limit_style_stop_order(client_order_id, order);
}
OrderType::StopMarket | OrderType::MarketIfTouched | OrderType::TrailingStopMarket => {
self.fill_market_order(client_order_id);
}
_ => {
log::error!(
"Cannot trigger stop order: invalid order type {}",
order.order_type()
);
}
}
}
fn trigger_limit_style_stop_order(&mut self, client_order_id: ClientOrderId, order: OrderAny) {
if order.is_triggered().is_some_and(|triggered| triggered) {
let liquidity_side = match (order.price(), order.trigger_price()) {
(Some(price), Some(trigger_price)) => Self::determine_triggered_limit_liquidity(
order.order_side(),
price,
trigger_price,
),
_ => LiquiditySide::Maker,
};
if let Some(mut cached_order) = self.cache.borrow_mut().order_mut(&client_order_id)
&& !matches!(
cached_order.liquidity_side(),
Some(LiquiditySide::Maker | LiquiditySide::Taker)
)
{
cached_order.set_liquidity_side(liquidity_side);
}
self.fill_limit_order(client_order_id);
return;
}
let event = self.create_order_triggered(&order);
let order = match self.cache.borrow_mut().update_order(&event) {
Ok(order) => order,
Err(e) => {
log::debug!(
"Failed to apply triggered event for {} before fill: {e}",
order.client_order_id(),
);
order
}
};
let order = self.order_snapshot(client_order_id).unwrap_or(order);
self.dispatch_order_event(event);
let trigger_price = order
.trigger_price()
.expect("Limit-style stop order must have a trigger price");
let price = order
.price()
.expect("Limit-style stop order must have a price");
let maker_inside = match order.order_side() {
OrderSide::Buy => self
.core
.ask
.is_some_and(|ask| trigger_price > price && price > ask),
OrderSide::Sell => self
.core
.bid
.is_some_and(|bid| trigger_price < price && price < bid),
};
if maker_inside {
if let Some(mut cached_order) = self.cache.borrow_mut().order_mut(&client_order_id) {
cached_order.set_liquidity_side(LiquiditySide::Maker);
}
self.resync_core_entry(client_order_id);
self.fill_limit_order(client_order_id);
return;
}
if self.core.is_limit_matched(order.order_side(), price) {
if order.is_post_only() {
self.delete_core_order(client_order_id);
self.cached_filled_qty.swap_remove(&client_order_id);
let event = self.create_order_rejected(
&order,
format!(
"POST_ONLY {} {} order limit px of {} would have been a TAKER: bid={}, ask={}",
order.order_type(),
order.order_side(),
price,
self.core
.bid
.map_or_else(|| "None".to_string(), |p| p.to_string()),
self.core
.ask
.map_or_else(|| "None".to_string(), |p| p.to_string())
)
.into(),
);
if let Err(e) = self.cache.borrow_mut().update_order(&event) {
log::debug!(
"Failed to apply rejected event for {} after post-only trigger: {e}",
order.client_order_id(),
);
}
self.dispatch_order_event(event);
return;
}
if let Some(mut cached_order) = self.cache.borrow_mut().order_mut(&client_order_id) {
cached_order.set_liquidity_side(LiquiditySide::Taker);
}
self.resync_core_entry(client_order_id);
self.fill_limit_order(client_order_id);
return;
}
if let Some(mut cached_order) = self.cache.borrow_mut().order_mut(&client_order_id) {
cached_order.set_liquidity_side(Self::determine_triggered_limit_liquidity(
order.order_side(),
price,
trigger_price,
));
}
self.resync_core_entry(client_order_id);
}
fn determine_triggered_limit_liquidity(
side: OrderSide,
price: Price,
trigger_price: Price,
) -> LiquiditySide {
if (side == OrderSide::Buy && trigger_price > price)
|| (side == OrderSide::Sell && trigger_price < price)
{
LiquiditySide::Maker
} else {
LiquiditySide::Taker
}
}
fn update_contingent_order(&mut self, order: &OrderAny, parent_quantity: Quantity) {
log::debug!(
"Updating contingent orders from {}",
order.client_order_id()
);
if let Some(linked_order_ids) = order.linked_order_ids() {
let parent_filled_qty = self
.cached_filled_qty
.get(&order.client_order_id())
.copied()
.unwrap_or(order.filled_qty());
let parent_leaves_qty = parent_quantity.saturating_sub(parent_filled_qty);
for client_order_id in linked_order_ids {
let child_order = match self.order_snapshot(*client_order_id) {
Some(order) => order,
None => panic!("Order {client_order_id} not found in cache."),
};
if child_order.is_active_local() {
continue;
}
let child_filled_qty = self
.cached_filled_qty
.get(&child_order.client_order_id())
.copied()
.unwrap_or(child_order.filled_qty());
if parent_leaves_qty.is_zero() {
self.cancel_order(&child_order, Some(false));
} else if child_filled_qty >= parent_leaves_qty {
self.cancel_order(&child_order, Some(false));
} else {
let child_leaves_qty = child_order.quantity().saturating_sub(child_filled_qty);
if child_leaves_qty != parent_leaves_qty {
let price = child_order.price();
let trigger_price = child_order.trigger_price();
self.update_order(
&child_order,
Some(parent_leaves_qty),
price,
trigger_price,
Some(false),
);
}
}
}
}
}
fn cancel_contingent_orders(&mut self, order: &OrderAny, excluded: &[ClientOrderId]) {
if let Some(linked_order_ids) = order.linked_order_ids() {
for client_order_id in linked_order_ids {
if excluded.contains(client_order_id) {
continue;
}
let contingent_order = match self.order_snapshot(*client_order_id) {
Some(order) => order,
None => panic!("Cannot find contingent order for {client_order_id}"),
};
if contingent_order.is_active_local() {
continue;
}
if !contingent_order.is_closed() {
self.cancel_order(&contingent_order, Some(false));
}
}
}
}
fn generate_order_submitted(&self, order: &OrderAny, account_id: AccountId) {
let ts_now = self.clock.borrow().timestamp_ns();
let event = OrderEventAny::Submitted(OrderSubmitted::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
account_id,
UUID4::new(),
ts_now,
ts_now,
));
self.dispatch_order_event(event);
}
fn create_order_rejected(&self, order: &OrderAny, reason: Ustr) -> OrderEventAny {
let ts_now = self.clock.borrow().timestamp_ns();
let account_id = order
.account_id()
.unwrap_or(self.account_ids.get(&order.trader_id()).unwrap().to_owned());
let due_post_only = reason.starts_with("POST_ONLY");
OrderEventAny::Rejected(OrderRejected::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
account_id,
reason,
UUID4::new(),
ts_now,
ts_now,
false,
due_post_only,
))
}
fn generate_order_rejected(&self, order: &OrderAny, reason: Ustr) {
let event = self.create_order_rejected(order, reason);
self.dispatch_order_event(event);
}
fn publish_order_initialized(&self, order: &OrderAny) {
let event = OrderEventAny::Initialized(order.init_event().clone());
msgbus::publish_order_event(
format!("events.order.{}", order.strategy_id()).into(),
&event,
);
}
fn create_order_accepted(
&self,
order: &OrderAny,
venue_order_id: VenueOrderId,
) -> OrderEventAny {
let ts_now = self.clock.borrow().timestamp_ns();
let account_id = order
.account_id()
.unwrap_or(self.account_ids.get(&order.trader_id()).unwrap().to_owned());
OrderEventAny::Accepted(OrderAccepted::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
venue_order_id,
account_id,
UUID4::new(),
ts_now,
ts_now,
false,
))
}
fn generate_order_accepted(&self, order: &OrderAny, venue_order_id: VenueOrderId) {
let event = self.create_order_accepted(order, venue_order_id);
self.dispatch_order_event(event);
}
#[expect(clippy::too_many_arguments)]
fn generate_order_modify_rejected(
&self,
trader_id: TraderId,
strategy_id: StrategyId,
instrument_id: InstrumentId,
client_order_id: ClientOrderId,
reason: Ustr,
venue_order_id: Option<VenueOrderId>,
account_id: Option<AccountId>,
) {
let ts_now = self.clock.borrow().timestamp_ns();
let event = OrderEventAny::ModifyRejected(OrderModifyRejected::new(
trader_id,
strategy_id,
instrument_id,
client_order_id,
reason,
UUID4::new(),
ts_now,
ts_now,
false,
venue_order_id,
account_id,
));
self.dispatch_order_event(event);
}
#[expect(clippy::too_many_arguments)]
fn generate_order_cancel_rejected(
&self,
trader_id: TraderId,
strategy_id: StrategyId,
account_id: AccountId,
instrument_id: InstrumentId,
client_order_id: ClientOrderId,
venue_order_id: Option<VenueOrderId>,
reason: Ustr,
) {
let ts_now = self.clock.borrow().timestamp_ns();
let event = OrderEventAny::CancelRejected(OrderCancelRejected::new(
trader_id,
strategy_id,
instrument_id,
client_order_id,
reason,
UUID4::new(),
ts_now,
ts_now,
false,
venue_order_id,
Some(account_id),
));
self.dispatch_order_event(event);
}
fn generate_order_updated(
&self,
order: &OrderAny,
quantity: Quantity,
price: Option<Price>,
trigger_price: Option<Price>,
protection_price: Option<Price>,
) {
let ts_now = self.clock.borrow().timestamp_ns();
let event = OrderUpdated::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
quantity,
UUID4::new(),
ts_now,
ts_now,
false,
order.venue_order_id(),
order.account_id(),
price,
trigger_price,
protection_price,
order.is_quote_quantity(),
);
self.pending_order_updates
.borrow_mut()
.entry(order.client_order_id())
.or_default()
.push(event);
self.dispatch_order_event(OrderEventAny::Updated(event));
}
fn generate_order_canceled(&self, order: &OrderAny, venue_order_id: VenueOrderId) {
let ts_now = self.clock.borrow().timestamp_ns();
let event = OrderEventAny::Canceled(OrderCanceled::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
UUID4::new(),
ts_now,
ts_now,
false,
Some(venue_order_id),
order.account_id(),
None,
));
self.dispatch_order_event(event);
}
fn create_order_triggered(&self, order: &OrderAny) -> OrderEventAny {
let ts_now = self.clock.borrow().timestamp_ns();
OrderEventAny::Triggered(OrderTriggered::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
UUID4::new(),
ts_now,
ts_now,
false,
order.venue_order_id(),
order.account_id(),
))
}
fn generate_order_expired(&self, order: &OrderAny) {
let ts_now = self.clock.borrow().timestamp_ns();
let event = OrderEventAny::Expired(OrderExpired::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
UUID4::new(),
ts_now,
ts_now,
false,
order.venue_order_id(),
order.account_id(),
));
self.dispatch_order_event(event);
}
#[expect(clippy::too_many_arguments)]
fn generate_order_filled(
&mut self,
order: &OrderAny,
venue_order_id: VenueOrderId,
venue_position_id: Option<PositionId>,
last_qty: Quantity,
last_px: Price,
quote_currency: Currency,
commission: Money,
liquidity_side: LiquiditySide,
) {
debug_assert!(
last_qty <= order.quantity(),
"Fill quantity {last_qty} exceeds order quantity {order_qty} for {client_order_id}",
order_qty = order.quantity(),
client_order_id = order.client_order_id()
);
let ts_now = self.clock.borrow().timestamp_ns();
let account_id = order
.account_id()
.unwrap_or(self.account_ids.get(&order.trader_id()).unwrap().to_owned());
let fill = OrderFilled::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
venue_order_id,
account_id,
self.ids_generator.generate_trade_id(ts_now),
order.order_side(),
order.order_type(),
last_qty,
last_px,
quote_currency,
liquidity_side,
UUID4::new(),
ts_now,
ts_now,
false,
venue_position_id,
Some(commission),
None,
);
self.record_pending_fill(&fill);
self.dispatch_order_event(OrderEventAny::Filled(fill));
}
fn record_pending_fill(&mut self, fill: &OrderFilled) {
if !self.config.use_reduce_only || self.instrument.is_spread() {
return;
}
self.purge_applied_fills();
let cache = self.cache.borrow();
let position_id = cache
.position_id(&fill.client_order_id)
.copied()
.or(fill.position_id)
.or_else(|| {
(self.oms_type == OmsType::Netting).then(|| {
PositionId::new(format!("{}-{}", fill.instrument_id, fill.strategy_id))
})
});
let opening_trade_id = position_id.and_then(|id| {
cache
.position(&id)
.and_then(|position| position.events.first().map(|event| event.trade_id))
});
let mut quantity_change = if fill.order_side == OrderSide::Buy {
fill.last_qty.as_decimal()
} else {
-fill.last_qty.as_decimal()
};
if matches!(self.instrument, InstrumentAny::CurrencyPair(_))
&& let Some(commission) = fill.commission
&& Some(commission.currency) == self.instrument.base_currency()
{
quantity_change -= commission.as_decimal();
}
self.pending_fills.insert(
fill.trade_id,
PendingFill {
client_order_id: fill.client_order_id,
position_id,
opening_trade_id,
quantity_change,
},
);
}
}
#[derive(Debug)]
struct PendingFill {
client_order_id: ClientOrderId,
position_id: Option<PositionId>,
opening_trade_id: Option<TradeId>,
quantity_change: Decimal,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ModifyOutcome {
Applied,
Rejected,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum OrderMatchMode {
All,
LastPriceStopTriggers,
}
#[derive(Debug)]
enum PostMatchOrderAction {
RemoveClosed,
Expire(OrderAny),
UpdateTrailing(OrderAny),
NoMaintenance,
}
fn post_match_order_action<F>(
order: &OrderAny,
support_gtd_orders: bool,
timestamp_ns: UnixNanos,
clone_order: F,
) -> PostMatchOrderAction
where
F: FnOnce(&OrderAny) -> OrderAny,
{
if order.is_closed() {
PostMatchOrderAction::RemoveClosed
} else if support_gtd_orders
&& order
.expire_time()
.is_some_and(|expire_ns| timestamp_ns >= expire_ns)
{
PostMatchOrderAction::Expire(clone_order(order))
} else if matches!(
order.order_type(),
OrderType::TrailingStopMarket | OrderType::TrailingStopLimit
) {
PostMatchOrderAction::UpdateTrailing(clone_order(order))
} else {
PostMatchOrderAction::NoMaintenance
}
}
fn write_filled_qty(order: &mut OrderAny, filled_qty: Quantity) {
match order {
OrderAny::Limit(o) => o.filled_qty = filled_qty,
OrderAny::LimitIfTouched(o) => o.filled_qty = filled_qty,
OrderAny::Market(o) => o.filled_qty = filled_qty,
OrderAny::MarketIfTouched(o) => o.filled_qty = filled_qty,
OrderAny::MarketToLimit(o) => o.filled_qty = filled_qty,
OrderAny::StopLimit(o) => o.filled_qty = filled_qty,
OrderAny::StopMarket(o) => o.filled_qty = filled_qty,
OrderAny::TrailingStopLimit(o) => o.filled_qty = filled_qty,
OrderAny::TrailingStopMarket(o) => o.filled_qty = filled_qty,
}
}
#[derive(Debug, Clone, Copy)]
struct BarTickSizes {
open: Quantity,
high: Quantity,
low: Quantity,
close: Quantity,
}
impl BarTickSizes {
fn from_volume(volume: Quantity, size_increment: Quantity) -> Self {
let precision_diff = FIXED_PRECISION.saturating_sub(volume.precision);
let scale = QuantityRaw::pow(10, u32::from(precision_diff));
let units = volume.raw() / scale;
let increment_units = (size_increment.raw() / scale).max(1);
let rounded_units = (units / increment_units) * increment_units;
let increments = rounded_units / increment_units;
let zero = Quantity::zero(volume.precision);
let size =
|increments| Quantity::from_raw(increments * increment_units * scale, volume.precision);
match increments {
0 => Self {
open: zero,
high: zero,
low: zero,
close: zero,
},
1 => Self {
open: zero,
high: zero,
low: zero,
close: size(1),
},
2 => Self {
open: zero,
high: size(1),
low: size(1),
close: zero,
},
3 => {
let path_size = size(1);
Self {
open: path_size,
high: path_size,
low: path_size,
close: zero,
}
}
_ => {
let path_increments = increments / 4;
let close_increments = increments - (path_increments * 3);
let path_size = size(path_increments);
Self {
open: path_size,
high: path_size,
low: path_size,
close: size(close_increments),
}
}
}
}
}
#[cfg(test)]
mod tests {
use std::{
cell::{Cell, RefCell},
collections::{HashMap, HashSet},
rc::Rc,
};
use nautilus_common::{
cache::Cache,
clock::TestClock,
messages::execution::{CancelAllOrders, ModifyOrder},
};
use nautilus_core::{UUID4, UnixNanos, correctness::CorrectnessError};
#[cfg(feature = "high-precision")]
use nautilus_model::orderbook::BookLevel;
use nautilus_model::{
data::{
Bar, BarType, DEPTH10_LEN, OrderBookDelta, OrderBookDeltas, OrderBookDepth10,
QuoteTick, TradeTick,
option_chain::OptionGreeks,
order::{BookOrder, OrderId},
},
enums::{
AccountType, AggressorSide, BookAction, BookType, ContingencyType, LiquiditySide,
OmsType, OrderSide, OrderStatus, OrderType, PositionSide, RecordFlag, TimeInForce,
TrailingOffsetType, TriggerType,
},
events::OrderEventAny,
identifiers::{AccountId, ClientOrderId, StrategyId, TradeId, TraderId, VenueOrderId},
instruments::{
Instrument, InstrumentAny,
stubs::{crypto_option_btc_deribit, crypto_perpetual_ethusdt, futures_contract_es},
},
orderbook::OrderBook,
orders::{Order, OrderAny, OrderTestBuilder, stubs::TestOrderEventStubs},
types::{Money, Price, Quantity, fixed::FIXED_PRECISION, quantity::QuantityRaw},
};
use proptest::prelude::*;
use rstest::rstest;
use rust_decimal::Decimal;
use super::{
BarTickSizes, OrderFilled, OrderMatchingEngine, Position, PositionId, PostMatchOrderAction,
post_match_order_action,
};
use crate::{
matching_engine::config::OrderMatchingEngineConfig,
models::{
fee::{FeeModel, FeeModelAny, FeeModelHandle},
fill::{FillModel, FillModelHandle},
},
};
fn assert_valid_bar_tick_sizes(volume: Quantity, size_increment: Quantity) {
let sizes = BarTickSizes::from_volume(volume, size_increment);
let total_raw = sizes.open.raw() + sizes.high.raw() + sizes.low.raw() + sizes.close.raw();
assert!(total_raw <= volume.raw());
for quantity in [sizes.open, sizes.high, sizes.low, sizes.close] {
assert_eq!(quantity.precision, volume.precision);
assert!(
OrderMatchingEngine::quantity_matches_precision(quantity, volume.precision),
"bar tick quantity {quantity} not aligned to precision {}",
volume.precision,
);
assert!(
size_increment.is_zero() || quantity.raw().is_multiple_of(size_increment.raw()),
"bar tick quantity {quantity} not aligned to increment {size_increment}",
);
}
if size_increment.is_positive() {
assert!(
volume.raw() - total_raw < size_increment.raw(),
"bar tick split left {} raw units from volume {volume} and increment {size_increment}",
volume.raw() - total_raw,
);
}
}
#[rstest]
#[case("100.009", "100.011", "100.000", true)]
#[case("100.009", "100.020", "100.008", false)]
#[case("100.010", "100.020", "100.000", false)]
fn test_bar_high_first_preserves_stored_distances(
#[case] open: &str,
#[case] high: &str,
#[case] low: &str,
#[case] expected: bool,
) {
let (mut engine, _, _) = collision_engine();
engine.config.bar_adaptive_high_low_ordering = true;
let mut prices = [Price::from(open), Price::from(high), Price::from(low)];
for price in &mut prices {
price.precision = 2;
}
let bar = Bar::new(
BarType::from("ETHUSDT-PERP.BINANCE-1-MINUTE-LAST-EXTERNAL"),
prices[0],
prices[1],
prices[2],
prices[0],
Quantity::from("1.000"),
1.into(),
1.into(),
);
assert_eq!(engine.bar_high_first(&bar), expected);
}
#[cfg(feature = "high-precision")]
#[rstest]
fn test_consume_trade_level_preserves_native_raw_units() {
let precision = if Quantity::from_raw_checked(0, 18).is_ok() {
18
} else {
FIXED_PRECISION
};
let size = Quantity::from_raw(2_000_000_000_000_000_000, precision);
let level = BookLevel::from_order(BookOrder::new(
OrderSide::Sell,
Price::from("1.00"),
size,
1,
));
let mut consumption = indexmap::IndexMap::default();
let mut remaining = size.raw();
OrderMatchingEngine::consume_trade_level(&mut consumption, &mut remaining, &level);
assert_eq!(remaining, 0);
assert_eq!(
consumption[&level.price.value.raw()],
(size.raw(), size.raw())
);
}
#[rstest]
fn test_post_match_order_action_does_not_clone_no_maintenance_order() {
let order = post_match_limit_order();
let clone_count = Cell::new(0);
let action = post_match_order_action(&order, true, UnixNanos::from(1_u64), |order| {
clone_count.set(clone_count.get() + 1);
order.clone()
});
assert!(matches!(action, PostMatchOrderAction::NoMaintenance));
assert_eq!(clone_count.get(), 0);
}
#[rstest]
#[case::spread(
InstrumentAny::FuturesSpread(nautilus_model::instruments::stubs::futures_spread_es()),
0
)]
#[case::outright(InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt()), 1)]
fn test_pending_fills_exclude_instruments_without_positions(
#[case] instrument: InstrumentAny,
#[case] expected_pending: usize,
) {
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
Default::default(),
);
let (order, fill) = pending_position_fill(
&instrument,
PositionId::from("POSITION-001"),
"OPEN",
OrderSide::Buy,
"1",
);
cache
.borrow_mut()
.add_order(order, None, None, false)
.unwrap();
engine.record_pending_fill(&fill);
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(fill))
.unwrap();
engine.purge_applied_fills();
assert_eq!(engine.pending_fills.len(), expected_pending);
}
#[rstest]
fn test_pending_fills_wait_for_position_acknowledgement(
#[values(OmsType::Netting, OmsType::Hedging)] oms_type: OmsType,
#[values(OrderSide::Buy, OrderSide::Sell)] closing_side: OrderSide,
) {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
oms_type,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
Default::default(),
);
let position_id = PositionId::from("POSITION-001");
let opening_side = if closing_side == OrderSide::Buy {
OrderSide::Sell
} else {
OrderSide::Buy
};
let (opening, opening_fill) =
pending_position_fill(&instrument, position_id, "OPEN", opening_side, "0.500");
let (closing, first_fill) = pending_position_fill(
&instrument,
position_id,
"CLOSE-FIRST",
closing_side,
"0.400",
);
let (_, second_fill) = pending_position_fill(
&instrument,
position_id,
"CLOSE-SECOND",
closing_side,
"0.100",
);
let (unrelated, unrelated_fill) = pending_position_fill(
&instrument,
PositionId::from("POSITION-002"),
"UNRELATED",
closing_side,
"0.200",
);
let position = Position::new(&instrument, opening_fill);
cache
.borrow_mut()
.add_order(opening, None, None, false)
.unwrap();
cache
.borrow_mut()
.add_position(&position, oms_type)
.unwrap();
cache
.borrow_mut()
.add_order(closing.clone(), Some(position_id), None, false)
.unwrap();
engine.record_pending_fill(&first_fill);
cache
.borrow_mut()
.add_order(unrelated, None, None, false)
.unwrap();
engine.record_pending_fill(&unrelated_fill);
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from("0.100")
);
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(first_fill.clone()))
.unwrap();
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from("0.100")
);
let position = cache
.borrow_mut()
.update_position_from_fill(position_id, &first_fill)
.unwrap();
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from("0.100")
);
assert!(!engine.pending_fills.contains_key(&first_fill.trade_id));
engine.record_pending_fill(&second_fill);
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from("0.000")
);
engine.reset();
assert!(engine.pending_fills.is_empty());
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from("0.100")
);
}
#[rstest]
#[case::base_fee("0.010 ETH", "0.89000")]
#[case::quote_fee("0.010 USDT", "0.90000")]
fn test_pending_spot_fills_include_base_currency_commission(
#[case] commission: &str,
#[case] expected: &str,
) {
let instrument = InstrumentAny::CurrencyPair(
nautilus_model::instruments::stubs::currency_pair_ethusdt(),
);
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Cash,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
Default::default(),
);
let position_id = PositionId::from("POSITION-001");
let (opening, opening_fill) =
pending_position_fill(&instrument, position_id, "OPEN", OrderSide::Buy, "0.50000");
let (_, mut increase_fill) = pending_position_fill(
&instrument,
position_id,
"INCREASE",
OrderSide::Buy,
"0.40000",
);
let (closing, _) = pending_position_fill(
&instrument,
position_id,
"CLOSE",
OrderSide::Sell,
"1.00000",
);
increase_fill.commission = Some(Money::from(commission));
let position = Position::new(&instrument, opening_fill);
cache
.borrow_mut()
.add_order(opening, None, None, false)
.unwrap();
cache
.borrow_mut()
.add_position(&position, OmsType::Netting)
.unwrap();
engine.record_pending_fill(&increase_fill);
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from(expected)
);
let position = cache
.borrow_mut()
.update_position_from_fill(position_id, &increase_fill)
.unwrap();
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from(expected)
);
assert!(engine.pending_fills.is_empty());
}
#[rstest]
fn test_pending_fills_survive_position_flip_and_archive_acknowledgement() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
Default::default(),
);
let position_id = PositionId::from("POSITION-001");
let (opening, opening_fill) =
pending_position_fill(&instrument, position_id, "OPEN", OrderSide::Buy, "10.000");
let (flipping, flip_fill) =
pending_position_fill(&instrument, position_id, "FLIP", OrderSide::Sell, "15.000");
let (closing, close_fill) =
pending_position_fill(&instrument, position_id, "CLOSE", OrderSide::Buy, "4.000");
for order in [opening, flipping, closing.clone()] {
cache
.borrow_mut()
.add_order(order, None, None, false)
.unwrap();
}
let mut position = Position::new(&instrument, opening_fill);
cache
.borrow_mut()
.add_position(&position, OmsType::Netting)
.unwrap();
engine.record_pending_fill(&flip_fill);
engine.record_pending_fill(&close_fill);
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from("1.000")
);
let (closing_flip, opening_flip) = flip_fill
.split_for_position_flip(Quantity::from("10.000"), Some(position_id), UUID4::new())
.unwrap();
position.apply(&closing_flip);
cache.borrow_mut().snapshot_position(&position).unwrap();
let position = Position::new(&instrument, opening_flip);
cache
.borrow_mut()
.add_position(&position, OmsType::Netting)
.unwrap();
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from("1.000")
);
assert!(!engine.pending_fills.contains_key(&flip_fill.trade_id));
assert!(engine.pending_fills.contains_key(&close_fill.trade_id));
let position = cache
.borrow_mut()
.update_position_from_fill(position_id, &close_fill)
.unwrap();
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from("1.000")
);
assert!(engine.pending_fills.is_empty());
let (_, flatten_fill) =
pending_position_fill(&instrument, position_id, "FLATTEN", OrderSide::Buy, "1.000");
let (_, reopen_fill) =
pending_position_fill(&instrument, position_id, "REOPEN", OrderSide::Sell, "3.000");
engine.record_pending_fill(&flatten_fill);
engine.record_pending_fill(&reopen_fill);
cache
.borrow_mut()
.update_position_from_fill(position_id, &flatten_fill)
.unwrap();
let closed = cache.borrow().position(&position_id).unwrap().clone();
cache.borrow_mut().snapshot_position(&closed).unwrap();
let position = Position::new(&instrument, reopen_fill);
cache
.borrow_mut()
.add_position_without_order(&position, OmsType::Netting)
.unwrap();
assert_eq!(
engine
.position_quantity_remaining(&closing, &position)
.unwrap(),
Quantity::from("3.000")
);
assert!(engine.pending_fills.is_empty());
}
#[rstest]
fn test_position_fills_sync_reduce_only_orders(
#[values(OrderSide::Buy, OrderSide::Sell)] opening_side: OrderSide,
#[values(OmsType::Netting, OmsType::Hedging)] oms_type: OmsType,
#[values(false, true)] deferred: bool,
#[values(OrderType::Limit, OrderType::StopMarket, OrderType::StopLimit)]
resting_type: OrderType,
#[values(false, true)] support_contingent_orders: bool,
#[values(false, true)] indexed: bool,
) {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L2_MBP,
oms_type,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
OrderMatchingEngineConfig {
support_contingent_orders,
..Default::default()
},
);
let position_id = PositionId::from("SYNC-POSITION");
let closing_side = if opening_side == OrderSide::Buy {
OrderSide::Sell
} else {
OrderSide::Buy
};
let (opening, mut opening_fill) = pending_position_fill(
&instrument,
position_id,
"SYNC-OPEN",
opening_side,
if support_contingent_orders {
"3.000"
} else {
"10.000"
},
);
let position_id = if oms_type == OmsType::Netting {
PositionId::new(format!("{}-{}", instrument.id(), opening.strategy_id()))
} else {
position_id
};
opening_fill.position_id = Some(position_id);
let mut position = Position::new(&instrument, opening_fill.clone());
cache
.borrow_mut()
.add_order(opening, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(opening_fill))
.unwrap();
for (id, qty) in [("SYNC-PARENT-A", "2.000"), ("SYNC-PARENT-B", "5.000")] {
let (parent, mut fill) =
pending_position_fill(&instrument, position_id, id, opening_side, qty);
fill.venue_order_id = VenueOrderId::from(id);
if support_contingent_orders {
position.apply(&fill);
}
cache
.borrow_mut()
.add_order(parent, Some(position_id), None, false)
.unwrap();
if support_contingent_orders {
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(fill))
.unwrap();
}
}
cache
.borrow_mut()
.add_position(&position, oms_type)
.unwrap();
engine
.account_ids
.insert(position.trader_id, position.account_id);
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = events.clone();
let handler_cache = cache.clone();
engine.set_event_handler(Rc::new(move |event| {
if !deferred || matches!(event, OrderEventAny::Accepted(_)) {
handler_cache.borrow_mut().update_order(&event).unwrap();
if let OrderEventAny::Filled(fill) = &event {
handler_cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
events_handler.borrow_mut().push(event);
}));
for (id, parent, reduce_only, assigned_position) in [
("SYNC-A", Some("SYNC-PARENT-A"), true, position_id),
("SYNC-B", Some("SYNC-PARENT-B"), true, position_id),
("SYNC-STANDALONE", None, true, position_id),
("SYNC-NON-REDUCE", None, false, position_id),
(
"SYNC-UNRELATED",
None,
true,
PositionId::from("OTHER-POSITION"),
),
] {
let mut builder = OrderTestBuilder::new(resting_type);
builder
.instrument_id(instrument.id())
.client_order_id(ClientOrderId::from(id))
.side(closing_side)
.quantity(Quantity::from("10.000"))
.reduce_only(reduce_only)
.submit(true);
if resting_type != OrderType::StopMarket {
builder.price(Price::from("2000.00"));
}
if resting_type != OrderType::Limit {
builder.trigger_price(Price::from("3000.00"));
}
if let Some(parent) = parent {
builder.parent_order_id(ClientOrderId::from(parent));
}
let mut order = builder.build();
cache
.borrow_mut()
.add_order(
order.clone(),
if !indexed && assigned_position == position_id {
None
} else {
Some(assigned_position)
},
None,
false,
)
.unwrap();
engine.accept_order(&mut order);
}
let (closing, _) = pending_position_fill(
&instrument,
position_id,
"SYNC-CLOSE",
closing_side,
"10.000",
);
cache
.borrow_mut()
.add_order(closing.clone(), Some(position_id), None, false)
.unwrap();
events.borrow_mut().clear();
for (quantity, expected_updates, expected_cancels) in [
(
"4.000",
vec![
(
"SYNC-A",
if support_contingent_orders {
"2.000"
} else {
"6.000"
},
),
(
"SYNC-B",
if support_contingent_orders {
"5.000"
} else {
"6.000"
},
),
("SYNC-STANDALONE", "6.000"),
],
Vec::new(),
),
(
"2.000",
if support_contingent_orders {
vec![("SYNC-B", "4.000"), ("SYNC-STANDALONE", "4.000")]
} else {
vec![
("SYNC-A", "4.000"),
("SYNC-B", "4.000"),
("SYNC-STANDALONE", "4.000"),
]
},
Vec::new(),
),
(
"4.000",
Vec::new(),
vec!["SYNC-A", "SYNC-B", "SYNC-STANDALONE"],
),
] {
let start = events.borrow().len();
engine
.apply_fills(
&closing,
&[(Price::from("1000.00"), Quantity::from(quantity))],
LiquiditySide::Taker,
Some(position_id),
Some(&position),
None,
)
.unwrap();
let events = events.borrow();
let emitted = &events[start..];
assert_eq!(
emitted.len(),
1 + expected_updates.len() + expected_cancels.len()
);
let OrderEventAny::Filled(fill) = &emitted[0] else {
panic!("Expected closing fill first")
};
assert_eq!(fill.client_order_id, closing.client_order_id());
assert_eq!(fill.last_qty, Quantity::from(quantity));
assert_eq!(fill.last_px, Price::from("1000.00"));
let mut updates = Vec::new();
let mut cancels = Vec::new();
for event in &emitted[1..] {
match event {
OrderEventAny::Updated(update) => {
assert_eq!(
update.price,
(resting_type != OrderType::StopMarket).then(|| Price::from("2000.00"))
);
assert_eq!(
update.trigger_price,
(resting_type != OrderType::Limit).then(|| Price::from("3000.00"))
);
updates.push((update.client_order_id.to_string(), update.quantity));
}
OrderEventAny::Canceled(cancel) => {
cancels.push(cancel.client_order_id.to_string());
}
other => panic!("Unexpected event {other:?}"),
}
}
updates.sort_by(|a, b| a.0.cmp(&b.0));
cancels.sort();
assert_eq!(
updates,
expected_updates
.into_iter()
.map(|(id, qty)| (id.to_string(), Quantity::from(qty)))
.collect::<Vec<_>>()
);
assert_eq!(cancels, expected_cancels);
}
if deferred {
for event in events.borrow().iter() {
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
}
let cache = cache.borrow();
assert_eq!(
cache.position(&position_id).unwrap().quantity,
Quantity::from("0.000")
);
for (id, quantity) in [
(
"SYNC-A",
if support_contingent_orders {
"2.000"
} else {
"4.000"
},
),
("SYNC-B", "4.000"),
("SYNC-STANDALONE", "4.000"),
] {
let id = ClientOrderId::from(id);
let order = cache.order(&id).unwrap();
assert_eq!(order.status(), OrderStatus::Canceled);
assert_eq!(order.quantity(), Quantity::from(quantity));
assert!(!engine.order_exists(id));
}
for id in ["SYNC-NON-REDUCE", "SYNC-UNRELATED"] {
let id = ClientOrderId::from(id);
let order = cache.order(&id).unwrap();
assert_eq!(order.status(), OrderStatus::Accepted);
assert_eq!(order.quantity(), Quantity::from("10.000"));
assert!(engine.order_exists(id));
}
}
#[rstest]
#[case(None, "7.000", "11.000", OrderStatus::PartiallyFilled)]
#[case(None, "9.000", "9.000", OrderStatus::PartiallyFilled)]
#[case(None, "10.000", "10.000", OrderStatus::Canceled)]
#[case(Some("10.000"), "7.000", "10.000", OrderStatus::PartiallyFilled)]
#[case(Some("9.000"), "7.000", "9.000", OrderStatus::PartiallyFilled)]
#[case(Some("8.000"), "7.000", "8.000", OrderStatus::Canceled)]
fn test_position_sync_accounts_for_prior_fills(
#[case] parent_filled: Option<&str>,
#[case] closing_quantity: &str,
#[case] expected_quantity: &str,
#[case] expected_status: OrderStatus,
#[values(false, true)] deferred: bool,
#[values(false, true)] use_reduce_only: bool,
) {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let position_id = PositionId::from("FLOOR-POSITION");
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L2_MBP,
OmsType::Hedging,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
OrderMatchingEngineConfig {
use_reduce_only,
..Default::default()
},
);
let opening_quantity = parent_filled.map_or(Quantity::from("18.000"), |quantity| {
Quantity::from("18.000") - Quantity::from(quantity)
});
let (opening, opening_fill) = pending_position_fill(
&instrument,
position_id,
"FLOOR-OPEN",
OrderSide::Buy,
&opening_quantity.to_string(),
);
let mut position = Position::new(&instrument, opening_fill.clone());
engine
.account_ids
.insert(position.trader_id, position.account_id);
cache
.borrow_mut()
.add_order(opening, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(opening_fill))
.unwrap();
let parent_id = parent_filled.map(|quantity| {
let (parent, mut fill) = pending_position_fill(
&instrument,
position_id,
"FLOOR-PARENT",
OrderSide::Buy,
quantity,
);
fill.venue_order_id = VenueOrderId::from("FLOOR-PARENT");
position.apply(&fill);
let parent_id = parent.client_order_id();
cache
.borrow_mut()
.add_order(parent, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(fill))
.unwrap();
parent_id
});
let mut builder = OrderTestBuilder::new(OrderType::Limit);
if let Some(parent_id) = parent_id {
builder.parent_order_id(parent_id);
}
let mut resting = builder
.instrument_id(instrument.id())
.client_order_id(ClientOrderId::from("FLOOR-RESTING"))
.side(OrderSide::Sell)
.quantity(Quantity::from("10.000"))
.price(Price::from("2000.00"))
.reduce_only(true)
.submit(true)
.build();
cache
.borrow_mut()
.add_order(resting.clone(), Some(position_id), None, false)
.unwrap();
let handler_cache = cache.clone();
engine.set_event_handler(Rc::new(move |event| {
handler_cache.borrow_mut().update_order(&event).unwrap();
}));
engine.accept_order(&mut resting);
let (_, mut prior_fill) = pending_position_fill(
&instrument,
position_id,
"FLOOR-RESTING",
OrderSide::Sell,
"8.000",
);
prior_fill.venue_order_id = resting.venue_order_id().unwrap();
prior_fill.order_type = OrderType::Limit;
position.apply(&prior_fill);
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(prior_fill))
.unwrap();
cache
.borrow_mut()
.add_position(&position, OmsType::Hedging)
.unwrap();
let (closing, _) = pending_position_fill(
&instrument,
position_id,
"FLOOR-CLOSE",
OrderSide::Sell,
closing_quantity,
);
cache
.borrow_mut()
.add_order(closing.clone(), Some(position_id), None, false)
.unwrap();
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = events.clone();
let handler_cache = cache.clone();
engine.set_event_handler(Rc::new(move |event| {
if !deferred {
handler_cache.borrow_mut().update_order(&event).unwrap();
if let OrderEventAny::Filled(fill) = &event {
handler_cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
events_handler.borrow_mut().push(event);
}));
engine
.apply_fills(
&closing,
&[(Price::from("1000.00"), Quantity::from(closing_quantity))],
LiquiditySide::Taker,
Some(position_id),
Some(&position),
None,
)
.unwrap();
let expected_quantity = Quantity::from(if use_reduce_only {
expected_quantity
} else {
"10.000"
});
let expected_status = if use_reduce_only {
expected_status
} else {
OrderStatus::PartiallyFilled
};
let updated = expected_quantity != Quantity::from("10.000");
let canceled = expected_status == OrderStatus::Canceled;
let events = events.borrow();
assert_eq!(
events.len(),
1 + usize::from(updated) + usize::from(canceled)
);
assert!(
matches!(&events[0], OrderEventAny::Filled(fill) if fill.last_qty == Quantity::from(closing_quantity))
);
if updated {
let OrderEventAny::Updated(update) = &events[1] else {
panic!("Expected remaining quantity update")
};
assert_eq!(update.client_order_id, resting.client_order_id());
assert_eq!(update.quantity, expected_quantity);
assert_eq!(update.price, Some(Price::from("2000.00")));
assert_eq!(update.trigger_price, None);
}
if canceled {
let OrderEventAny::Canceled(cancel) = events.last().unwrap() else {
panic!("Expected cancellation with no remaining capacity")
};
assert_eq!(cancel.client_order_id, resting.client_order_id());
}
if deferred {
for event in events.iter() {
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
}
let cache = cache.borrow();
let resting = cache.order(&resting.client_order_id()).unwrap();
assert_eq!(resting.filled_qty(), Quantity::from("8.000"));
assert_eq!(resting.quantity(), expected_quantity);
assert_eq!(
resting.leaves_qty(),
expected_quantity - Quantity::from("8.000")
);
assert_eq!(resting.status(), expected_status);
assert_eq!(engine.order_exists(resting.client_order_id()), !canceled);
assert_eq!(
cache.position(&position_id).unwrap().quantity,
Quantity::from("10.000") - Quantity::from(closing_quantity)
);
}
#[rstest]
#[case(("0.000", "0.000"), (None, None), "open", (["6.000", "4.000"], [Some("6.000"), Some("4.000")]), false)]
#[case(("2.000", "3.000"), (None, None), "open", (["8.000", "6.000"], [Some("9.000"), Some("7.000")]), false)]
#[case(("2.000", "3.000"), (Some("7.000"), Some("6.000")), "open", (["7.000", "6.000"], [Some("6.000"), None]), false)]
#[case(("2.000", "3.000"), (Some("2.000"), None), "open", (["2.000", "2.000"], [None, None]), true)]
#[case(("2.000", "3.000"), (None, Some("3.000")), "open", (["8.000", "6.000"], [Some("3.000"), None]), true)]
#[case(("0.000", "0.000"), (None, None), "closed", (["6.000", "4.000"], [None, None]), false)]
#[case(("0.000", "0.000"), (None, None), "local", (["6.000", "4.000"], [None, None]), false)]
#[case(("0.000", "0.000"), (None, None), "cancellation_unacknowledged", (["6.000", "4.000"], [None, None]), false)]
fn test_position_sync_resizes_mixed_ouo_sibling(
#[case] filled: (&str, &str),
#[case] parents: (Option<&str>, Option<&str>),
#[case] sibling_state: &str,
#[case] expected: ([&str; 2], [Option<&str>; 2]),
#[case] first_cancel: bool,
#[values(0, 1, 2)] delivery: usize,
#[values(false, true)] support_contingent_orders: bool,
) {
let (source_filled, sibling_filled) = filled;
let (source_parent, sibling_parent) = parents;
let (source_quantities, sibling_updates) = expected;
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let position_id = PositionId::from("MIXED-POSITION");
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L2_MBP,
OmsType::Hedging,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
OrderMatchingEngineConfig {
support_contingent_orders,
..Default::default()
},
);
let opening_quantity = Quantity::from("10.000")
+ Quantity::from(source_filled)
+ Quantity::from(sibling_filled);
let (opening, opening_fill) = pending_position_fill(
&instrument,
position_id,
"MIXED-OPEN",
OrderSide::Buy,
&opening_quantity.to_string(),
);
let mut position = Position::new(&instrument, opening_fill.clone());
engine
.account_ids
.insert(position.trader_id, position.account_id);
cache
.borrow_mut()
.add_order(opening, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(opening_fill))
.unwrap();
let handler_cache = cache.clone();
engine.set_event_handler(Rc::new(move |event| {
handler_cache.borrow_mut().update_order(&event).unwrap();
}));
for (id, sibling, reduce_only, filled, parent_quantity) in [
("MIXED-A", "MIXED-B", true, source_filled, source_parent),
("MIXED-B", "MIXED-A", false, sibling_filled, sibling_parent),
] {
let mut builder = OrderTestBuilder::new(OrderType::Limit);
if let Some(quantity) = parent_quantity {
let parent_id = format!("{id}-PARENT");
let (parent, mut fill) = pending_position_fill(
&instrument,
position_id,
&parent_id,
OrderSide::Buy,
quantity,
);
fill.venue_order_id = VenueOrderId::from(parent_id.as_str());
cache
.borrow_mut()
.add_order(parent, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(fill))
.unwrap();
builder.parent_order_id(ClientOrderId::from(parent_id));
}
let mut order = builder
.instrument_id(instrument.id())
.client_order_id(ClientOrderId::from(id))
.side(OrderSide::Sell)
.quantity(Quantity::from("10.000"))
.price(Price::from("2000.00"))
.reduce_only(reduce_only)
.contingency_type(ContingencyType::Ouo)
.linked_order_ids(vec![ClientOrderId::from(sibling)])
.submit(sibling_state != "local" || reduce_only)
.build();
cache
.borrow_mut()
.add_order(order.clone(), Some(position_id), None, false)
.unwrap();
if sibling_state != "local" || reduce_only {
engine.accept_order(&mut order);
}
if Quantity::from(filled).non_zero() {
let (_, mut fill) =
pending_position_fill(&instrument, position_id, id, OrderSide::Sell, filled);
fill.venue_order_id = order.venue_order_id().unwrap();
fill.order_type = OrderType::Limit;
position.apply(&fill);
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(fill))
.unwrap();
}
if !reduce_only && sibling_state == "closed" {
engine.cancel_order(&order, Some(false));
}
}
cache
.borrow_mut()
.add_position(&position, OmsType::Hedging)
.unwrap();
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = events.clone();
let handler_cache = cache.clone();
engine.set_event_handler(Rc::new(move |event| {
if delivery == 0 {
handler_cache.borrow_mut().update_order(&event).unwrap();
if let OrderEventAny::Filled(fill) = &event {
handler_cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
events_handler.borrow_mut().push(event);
}));
if sibling_state == "cancellation_unacknowledged" {
let sibling = engine
.order_snapshot(ClientOrderId::from("MIXED-B"))
.unwrap();
engine.cancel_order(&sibling, Some(false));
}
let (closing, _) = pending_position_fill(
&instrument,
position_id,
"MIXED-CLOSE",
OrderSide::Sell,
"10.000",
);
cache
.borrow_mut()
.add_order(closing.clone(), Some(position_id), None, false)
.unwrap();
let mut acknowledged = 0;
let mut source_quantity = Quantity::from("10.000");
let mut sibling_quantity = Quantity::from("10.000");
let mut source_canceled = false;
let mut sibling_canceled =
matches!(sibling_state, "closed" | "cancellation_unacknowledged");
for (step, (quantity, remaining)) in
[("4.000", "6.000"), ("2.000", "4.000"), ("4.000", "0.000")]
.into_iter()
.enumerate()
{
let start = events.borrow().len();
engine
.apply_fills(
&closing,
&[(Price::from("1000.00"), Quantity::from(quantity))],
LiquiditySide::Taker,
Some(position_id),
Some(&position),
None,
)
.unwrap();
let mut expected = vec![("fill", "MIXED-CLOSE", Quantity::from(quantity))];
if !source_canceled {
if step < 2 {
let target = if support_contingent_orders {
Quantity::from(source_quantities[step])
} else {
Quantity::from(source_filled) + Quantity::from(remaining)
};
if target != source_quantity {
expected.push(("update", "MIXED-A", target));
source_quantity = target;
if support_contingent_orders && source_parent == Some(source_filled) {
expected.push(("cancel", "MIXED-A", Quantity::zero(3)));
source_canceled = true;
if !sibling_canceled && sibling_state != "local" {
expected.push(("cancel", "MIXED-B", Quantity::zero(3)));
sibling_canceled = true;
}
} else if support_contingent_orders {
if let Some(target) = sibling_updates[step] {
sibling_quantity = Quantity::from(target);
expected.push(("update", "MIXED-B", sibling_quantity));
}
if step == 0 && first_cancel {
expected.push(("cancel", "MIXED-B", Quantity::zero(3)));
sibling_canceled = true;
}
}
}
} else {
expected.push(("cancel", "MIXED-A", Quantity::zero(3)));
source_canceled = true;
if support_contingent_orders && !sibling_canceled && sibling_state != "local" {
expected.push(("cancel", "MIXED-B", Quantity::zero(3)));
sibling_canceled = true;
}
}
}
let recorded = events.borrow();
let actual: Vec<_> = recorded[start..]
.iter()
.map(|event| match event {
OrderEventAny::Filled(fill) => {
assert_eq!(fill.last_px, Price::from("1000.00"));
("fill", fill.client_order_id.as_str(), fill.last_qty)
}
OrderEventAny::Updated(update) => {
assert_eq!(update.price, Some(Price::from("2000.00")));
assert_eq!(update.trigger_price, None);
("update", update.client_order_id.as_str(), update.quantity)
}
OrderEventAny::Canceled(cancel) => {
("cancel", cancel.client_order_id.as_str(), Quantity::zero(3))
}
other => panic!("Unexpected event {other:?}"),
})
.collect();
assert_eq!(actual, expected);
drop(recorded);
if delivery == 2 {
let end = events.borrow().len() - 1;
for event in &events.borrow()[acknowledged..end] {
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
acknowledged = end;
}
let before = events.borrow().len();
let ids = engine.reduce_only_order_ids(position_id);
engine
.sync_reduce_only_orders(&closing, &position, &ids)
.unwrap();
assert_eq!(events.borrow().len(), before);
}
if delivery != 0 {
for event in &events.borrow()[acknowledged..] {
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
}
let cache = cache.borrow();
for (id, filled, quantity, canceled) in [
("MIXED-A", source_filled, source_quantity, source_canceled),
(
"MIXED-B",
sibling_filled,
sibling_quantity,
sibling_canceled,
),
] {
let order = cache.order(&ClientOrderId::from(id)).unwrap();
assert_eq!(order.quantity(), quantity);
assert_eq!(order.filled_qty(), Quantity::from(filled));
assert_eq!(order.leaves_qty(), quantity - Quantity::from(filled));
assert_eq!(
order.status(),
if canceled {
OrderStatus::Canceled
} else if sibling_state == "local" {
OrderStatus::Initialized
} else if Quantity::from(filled).is_zero() {
OrderStatus::Accepted
} else {
OrderStatus::PartiallyFilled
}
);
assert_eq!(
engine.order_exists(order.client_order_id()),
!canceled && sibling_state != "local"
);
}
assert_eq!(
cache.position(&position_id).unwrap().quantity,
Quantity::from("0.000")
);
assert_eq!(
cache.position(&position_id).unwrap().side,
PositionSide::Flat
);
}
#[rstest]
fn test_position_sync_does_not_resize_order_being_filled(
#[values(false, true)] deferred: bool,
) {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let position_id = PositionId::from("REENTRANT-POSITION");
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L2_MBP,
OmsType::Hedging,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
Default::default(),
);
let (opening, opening_fill) = pending_position_fill(
&instrument,
position_id,
"REENTRANT-OPEN",
OrderSide::Buy,
"6.000",
);
let position = Position::new(&instrument, opening_fill.clone());
engine
.account_ids
.insert(position.trader_id, position.account_id);
cache
.borrow_mut()
.add_order(opening, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(opening_fill))
.unwrap();
cache
.borrow_mut()
.add_position(&position, OmsType::Hedging)
.unwrap();
for (id, price, size) in [(1, "1000.00", "4.000"), (2, "999.00", "5.000")] {
engine
.process_order_book_delta(&OrderBookDelta::new(
instrument.id(),
BookAction::Add,
BookOrder::new(OrderSide::Buy, Price::from(price), Quantity::from(size), id),
0,
id,
UnixNanos::from(id),
UnixNanos::from(id),
))
.unwrap();
}
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = events.clone();
let handler_cache = cache.clone();
engine.set_event_handler(Rc::new(move |event| {
if !deferred || matches!(event, OrderEventAny::Accepted(_)) {
handler_cache.borrow_mut().update_order(&event).unwrap();
if let OrderEventAny::Filled(fill) = &event {
handler_cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
events_handler.borrow_mut().push(event);
}));
for (id, sibling, reduce_only) in [
("REENTRANT-A", "REENTRANT-B", true),
("REENTRANT-B", "REENTRANT-A", false),
] {
let mut builder = OrderTestBuilder::new(OrderType::Limit);
builder
.instrument_id(instrument.id())
.client_order_id(ClientOrderId::from(id))
.side(OrderSide::Sell)
.quantity(Quantity::from("10.000"))
.price(Price::from(if reduce_only { "2000.00" } else { "999.00" }))
.reduce_only(reduce_only)
.contingency_type(ContingencyType::Ouo)
.linked_order_ids(vec![ClientOrderId::from(sibling)])
.submit(true);
let mut order = builder.build();
order.set_liquidity_side(LiquiditySide::Taker);
cache
.borrow_mut()
.add_order(order.clone(), Some(position_id), None, false)
.unwrap();
engine.accept_order(&mut order);
}
events.borrow_mut().clear();
engine.iterate(UnixNanos::from(3), AggressorSide::NoAggressor);
if deferred {
for event in events.borrow().iter() {
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
}
let cache = cache.borrow();
let filled = cache.order(&ClientOrderId::from("REENTRANT-B")).unwrap();
assert_eq!(filled.quantity(), Quantity::from("10.000"));
assert_eq!(filled.filled_qty(), Quantity::from("9.000"));
assert_eq!(filled.leaves_qty(), Quantity::from("1.000"));
assert_eq!(filled.overfill_qty(), Quantity::from("0.000"));
assert_eq!(filled.status(), OrderStatus::PartiallyFilled);
assert_eq!(
cache.position(&position_id).unwrap().quantity,
Quantity::from("3.000")
);
assert_eq!(
cache.position(&position_id).unwrap().side,
PositionSide::Short
);
let recorded = events.borrow();
let actual: Vec<_> = recorded
.iter()
.map(|event| match event {
OrderEventAny::Filled(fill) => {
("fill", fill.client_order_id.as_str(), fill.last_qty)
}
OrderEventAny::Updated(update) => {
("update", update.client_order_id.as_str(), update.quantity)
}
OrderEventAny::Canceled(cancel) => {
("cancel", cancel.client_order_id.as_str(), Quantity::zero(3))
}
other => panic!("Unexpected event {other:?}"),
})
.collect();
assert_eq!(
actual,
vec![
("fill", "REENTRANT-B", Quantity::from("4.000")),
("update", "REENTRANT-A", Quantity::from("6.000")),
("update", "REENTRANT-A", Quantity::from("2.000")),
("fill", "REENTRANT-B", Quantity::from("5.000")),
("update", "REENTRANT-A", Quantity::from("1.000")),
("cancel", "REENTRANT-A", Quantity::zero(3)),
]
);
}
#[rstest]
#[case("5.000", "5.000", false)]
#[case("10.000", "0.000", true)]
fn test_position_sync_handles_unacknowledged_sibling_acceptance(
#[case] closing_quantity: &str,
#[case] remaining_quantity: &str,
#[case] canceled: bool,
#[values(false, true)] deferred: bool,
) {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let position_id = PositionId::from("REENTRANT-POSITION");
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L2_MBP,
OmsType::Hedging,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
Default::default(),
);
let (opening, opening_fill) = pending_position_fill(
&instrument,
position_id,
"REENTRANT-OPEN",
OrderSide::Buy,
"10.000",
);
let position = Position::new(&instrument, opening_fill.clone());
engine
.account_ids
.insert(position.trader_id, position.account_id);
cache
.borrow_mut()
.add_order(opening, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(opening_fill))
.unwrap();
cache
.borrow_mut()
.add_position(&position, OmsType::Hedging)
.unwrap();
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = events.clone();
let handler_cache = cache.clone();
let sibling_id = ClientOrderId::from("ACCEPT-B");
engine.set_event_handler(Rc::new(move |event| {
let id = match &event {
OrderEventAny::Accepted(event) => event.client_order_id,
OrderEventAny::Filled(event) => event.client_order_id,
OrderEventAny::Canceled(event) => event.client_order_id,
OrderEventAny::Updated(event) => event.client_order_id,
other => panic!("Unexpected event {other:?}"),
};
let applied =
id != sibling_id && (!deferred || matches!(event, OrderEventAny::Accepted(_)));
if applied {
handler_cache.borrow_mut().update_order(&event).unwrap();
if let OrderEventAny::Filled(fill) = &event {
handler_cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
events_handler.borrow_mut().push((event, applied));
}));
for (id, sibling, reduce_only) in [
("ACCEPT-A", "ACCEPT-B", true),
("ACCEPT-B", "ACCEPT-A", false),
] {
let mut order = OrderTestBuilder::new(OrderType::Limit)
.instrument_id(instrument.id())
.client_order_id(ClientOrderId::from(id))
.side(OrderSide::Sell)
.quantity(Quantity::from("10.000"))
.price(Price::from("2000.00"))
.reduce_only(reduce_only)
.contingency_type(ContingencyType::Ouo)
.linked_order_ids(vec![ClientOrderId::from(sibling)])
.submit(true)
.build();
cache
.borrow_mut()
.add_order(order.clone(), Some(position_id), None, false)
.unwrap();
engine.accept_order(&mut order);
}
assert_eq!(
cache.borrow().order(&sibling_id).unwrap().status(),
OrderStatus::Submitted
);
assert!(engine.order_exists(sibling_id));
let (closing, _) = pending_position_fill(
&instrument,
position_id,
"ACCEPT-CLOSE",
OrderSide::Sell,
closing_quantity,
);
cache
.borrow_mut()
.add_order(closing.clone(), Some(position_id), None, false)
.unwrap();
engine
.apply_fills(
&closing,
&[(Price::from("1000.00"), Quantity::from(closing_quantity))],
LiquiditySide::Taker,
Some(position_id),
Some(&position),
None,
)
.unwrap();
let ids = engine.reduce_only_order_ids(position_id);
engine
.sync_reduce_only_orders(&closing, &position, &ids)
.unwrap();
let events = events.borrow();
let actual: Vec<_> = events
.iter()
.map(|(event, _)| match event {
OrderEventAny::Accepted(event) => ("accepted", event.client_order_id.as_str()),
OrderEventAny::Filled(fill) => {
assert_eq!(fill.last_qty, Quantity::from(closing_quantity));
assert_eq!(fill.last_px, Price::from("1000.00"));
("filled", fill.client_order_id.as_str())
}
OrderEventAny::Updated(event) => {
assert_eq!(event.quantity, Quantity::from("5.000"));
assert_eq!(event.price, Some(Price::from("2000.00")));
assert_eq!(event.trigger_price, None);
("updated", event.client_order_id.as_str())
}
OrderEventAny::Canceled(event) => ("canceled", event.client_order_id.as_str()),
other => panic!("Unexpected event {other:?}"),
})
.collect();
let mut expected = vec![
("accepted", "ACCEPT-A"),
("accepted", "ACCEPT-B"),
("filled", "ACCEPT-CLOSE"),
];
if canceled {
expected.extend([("canceled", "ACCEPT-A"), ("canceled", "ACCEPT-B")]);
} else {
expected.push(("updated", "ACCEPT-A"));
}
assert_eq!(actual, expected);
assert_eq!(engine.order_exists(sibling_id), !canceled);
for (event, applied) in events.iter() {
if !applied {
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
}
let cache = cache.borrow();
for id in ["ACCEPT-A", "ACCEPT-B"] {
let order = cache.order(&ClientOrderId::from(id)).unwrap();
let quantity = Quantity::from(if !canceled && id == "ACCEPT-A" {
"5.000"
} else {
"10.000"
});
assert_eq!(
order.status(),
if canceled {
OrderStatus::Canceled
} else {
OrderStatus::Accepted
}
);
assert_eq!(order.quantity(), quantity);
assert_eq!(order.filled_qty(), Quantity::from("0.000"));
assert_eq!(order.leaves_qty(), quantity);
}
assert_eq!(
cache.position(&position_id).unwrap().quantity,
Quantity::from(remaining_quantity)
);
assert_eq!(
cache.position(&position_id).unwrap().side,
if canceled {
PositionSide::Flat
} else {
PositionSide::Long
}
);
}
#[rstest]
fn test_position_sync_mixed_ouo_does_not_match_recursively(#[values(0, 1, 2)] delivery: usize) {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let position_id = PositionId::from("REENTRANT-POSITION");
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L2_MBP,
OmsType::Hedging,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
Default::default(),
);
let (opening, opening_fill) = pending_position_fill(
&instrument,
position_id,
"REENTRANT-OPEN",
OrderSide::Buy,
"10.000",
);
let position = Position::new(&instrument, opening_fill.clone());
engine
.account_ids
.insert(position.trader_id, position.account_id);
cache
.borrow_mut()
.add_order(opening, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(opening_fill))
.unwrap();
cache
.borrow_mut()
.add_position(&position, OmsType::Hedging)
.unwrap();
for (id, price, size) in [(1, "1000.00", "1.000"), (2, "999.00", "9.000")] {
engine
.process_order_book_delta(&OrderBookDelta::new(
instrument.id(),
BookAction::Add,
BookOrder::new(OrderSide::Buy, Price::from(price), Quantity::from(size), id),
0,
id,
UnixNanos::from(id),
UnixNanos::from(id),
))
.unwrap();
}
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = events.clone();
let handler_cache = cache.clone();
engine.set_event_handler(Rc::new(move |event| {
if delivery == 0 || matches!(event, OrderEventAny::Accepted(_)) {
handler_cache.borrow_mut().update_order(&event).unwrap();
if let OrderEventAny::Filled(fill) = &event {
handler_cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
events_handler.borrow_mut().push(event);
}));
for (id, sibling, reduce_only) in [
("REENTRANT-A", "REENTRANT-B", true),
("REENTRANT-B", "REENTRANT-A", false),
] {
let mut builder = OrderTestBuilder::new(OrderType::Limit);
builder
.instrument_id(instrument.id())
.client_order_id(ClientOrderId::from(id))
.side(OrderSide::Sell)
.quantity(Quantity::from("10.000"))
.price(Price::from("999.00"))
.reduce_only(reduce_only)
.contingency_type(ContingencyType::Ouo)
.linked_order_ids(vec![ClientOrderId::from(sibling)])
.submit(true);
let mut order = builder.build();
order.set_liquidity_side(LiquiditySide::Taker);
cache
.borrow_mut()
.add_order(order.clone(), Some(position_id), None, false)
.unwrap();
engine.accept_order(&mut order);
}
events.borrow_mut().clear();
let (mut closing, _) = pending_position_fill(
&instrument,
position_id,
"REENTRANT-CLOSE",
OrderSide::Sell,
"4.000",
);
cache
.borrow_mut()
.add_order(closing.clone(), Some(position_id), None, false)
.unwrap();
engine.process_order(&mut closing, position.account_id);
let mut acknowledged = 0;
if delivery == 2 {
for event in &events.borrow()[..5] {
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
acknowledged = 5;
}
assert_eq!(
engine
.position_quantity_remaining(
&closing,
&cache.borrow().position(&position_id).unwrap()
)
.unwrap(),
Quantity::from("6.000")
);
for id in ["REENTRANT-A", "REENTRANT-B"] {
let order = engine.order_snapshot(ClientOrderId::from(id)).unwrap();
assert_eq!(order.quantity(), Quantity::from("6.000"));
assert_eq!(order.filled_qty(), Quantity::from("0.000"));
assert_eq!(order.leaves_qty(), Quantity::from("6.000"));
}
let (mut flattening, _) = pending_position_fill(
&instrument,
position_id,
"REENTRANT-FLAT",
OrderSide::Sell,
"6.000",
);
cache
.borrow_mut()
.add_order(flattening.clone(), Some(position_id), None, false)
.unwrap();
engine.process_order(&mut flattening, position.account_id);
let recorded = events.borrow();
let actual: Vec<_> = recorded
.iter()
.map(|event| match event {
OrderEventAny::Filled(fill) => (
"fill",
fill.client_order_id.as_str(),
fill.last_qty,
Some(fill.last_px),
),
OrderEventAny::Updated(update) => {
assert_eq!(update.trigger_price, None);
(
"update",
update.client_order_id.as_str(),
update.quantity,
update.price,
)
}
OrderEventAny::Canceled(cancel) => (
"cancel",
cancel.client_order_id.as_str(),
Quantity::zero(3),
None,
),
other => panic!("Unexpected event {other:?}"),
})
.collect();
assert_eq!(
actual,
vec![
(
"fill",
"REENTRANT-CLOSE",
Quantity::from("1.000"),
Some(Price::from("1000.00"))
),
(
"update",
"REENTRANT-A",
Quantity::from("9.000"),
Some(Price::from("999.00"))
),
(
"update",
"REENTRANT-B",
Quantity::from("9.000"),
Some(Price::from("999.00"))
),
(
"fill",
"REENTRANT-CLOSE",
Quantity::from("3.000"),
Some(Price::from("999.00"))
),
(
"update",
"REENTRANT-A",
Quantity::from("6.000"),
Some(Price::from("999.00"))
),
(
"update",
"REENTRANT-B",
Quantity::from("6.000"),
Some(Price::from("999.00"))
),
(
"fill",
"REENTRANT-FLAT",
Quantity::from("1.000"),
Some(Price::from("1000.00"))
),
(
"update",
"REENTRANT-A",
Quantity::from("5.000"),
Some(Price::from("999.00"))
),
(
"update",
"REENTRANT-B",
Quantity::from("5.000"),
Some(Price::from("999.00"))
),
(
"fill",
"REENTRANT-FLAT",
Quantity::from("5.000"),
Some(Price::from("999.00"))
),
("cancel", "REENTRANT-A", Quantity::zero(3), None),
("cancel", "REENTRANT-B", Quantity::zero(3), None),
]
);
if delivery != 0 {
for event in &recorded[acknowledged..] {
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
}
let cache = cache.borrow();
assert_eq!(
cache.position(&position_id).unwrap().quantity,
Quantity::from("0.000")
);
assert_eq!(
cache.position(&position_id).unwrap().side,
PositionSide::Flat
);
for id in ["REENTRANT-A", "REENTRANT-B"] {
let order = cache.order(&ClientOrderId::from(id)).unwrap();
assert_eq!(order.status(), OrderStatus::Canceled);
assert_eq!(order.quantity(), Quantity::from("5.000"));
assert_eq!(order.filled_qty(), Quantity::from("0.000"));
assert_eq!(order.leaves_qty(), Quantity::from("5.000"));
assert!(!engine.order_exists(order.client_order_id()));
}
}
#[rstest]
fn test_position_sync_does_not_match_recursively(#[values(false, true)] deferred: bool) {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let position_id = PositionId::from("REENTRANT-POSITION");
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L2_MBP,
OmsType::Hedging,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
Default::default(),
);
let (opening, opening_fill) = pending_position_fill(
&instrument,
position_id,
"REENTRANT-OPEN",
OrderSide::Buy,
"10.000",
);
let position = Position::new(&instrument, opening_fill.clone());
engine
.account_ids
.insert(position.trader_id, position.account_id);
cache
.borrow_mut()
.add_order(opening, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(opening_fill))
.unwrap();
cache
.borrow_mut()
.add_position(&position, OmsType::Hedging)
.unwrap();
let (parent, mut parent_fill) = pending_position_fill(
&instrument,
position_id,
"REENTRANT-PARENT",
OrderSide::Buy,
"2.000",
);
parent_fill.venue_order_id = VenueOrderId::from("REENTRANT-PARENT");
cache
.borrow_mut()
.add_order(parent, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(parent_fill))
.unwrap();
for (id, price, size) in [(1, "1000.00", "1.000"), (2, "999.00", "9.000")] {
engine
.process_order_book_delta(&OrderBookDelta::new(
instrument.id(),
BookAction::Add,
BookOrder::new(OrderSide::Buy, Price::from(price), Quantity::from(size), id),
0,
id,
UnixNanos::from(id),
UnixNanos::from(id),
))
.unwrap();
}
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = events.clone();
let handler_cache = cache.clone();
engine.set_event_handler(Rc::new(move |event| {
if !deferred || matches!(event, OrderEventAny::Accepted(_)) {
handler_cache.borrow_mut().update_order(&event).unwrap();
if let OrderEventAny::Filled(fill) = &event {
handler_cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
events_handler.borrow_mut().push(event);
}));
for (id, quantity, parent) in [
("REENTRANT-A", "2.000", Some("REENTRANT-PARENT")),
("REENTRANT-B", "10.000", None),
] {
let mut builder = OrderTestBuilder::new(OrderType::Limit);
builder
.instrument_id(instrument.id())
.client_order_id(ClientOrderId::from(id))
.side(OrderSide::Sell)
.quantity(Quantity::from(quantity))
.price(Price::from("999.00"))
.reduce_only(true)
.submit(true);
if let Some(parent) = parent {
builder.parent_order_id(ClientOrderId::from(parent));
}
let mut order = builder.build();
order.set_liquidity_side(LiquiditySide::Taker);
cache
.borrow_mut()
.add_order(order.clone(), Some(position_id), None, false)
.unwrap();
engine.accept_order(&mut order);
}
events.borrow_mut().clear();
assert_eq!(engine.core.iterate_asks().len(), 2);
assert_eq!(
cache.borrow().position(&position_id).unwrap().quantity,
Quantity::from("10.000")
);
engine.iterate(UnixNanos::from(3), AggressorSide::NoAggressor);
let events = events.borrow();
let fills: Vec<_> = events
.iter()
.filter_map(|event| match event {
OrderEventAny::Filled(fill) => Some((
fill.client_order_id.to_string(),
fill.last_qty,
fill.last_px,
)),
_ => None,
})
.collect();
assert_eq!(
fills,
vec![
(
"REENTRANT-A".to_string(),
Quantity::from("1.000"),
Price::from("1000.00")
),
(
"REENTRANT-A".to_string(),
Quantity::from("1.000"),
Price::from("999.00")
),
(
"REENTRANT-B".to_string(),
Quantity::from("1.000"),
Price::from("1000.00")
),
(
"REENTRANT-B".to_string(),
Quantity::from("7.000"),
Price::from("999.00")
),
]
);
assert!(
!events
.iter()
.any(|event| matches!(event, OrderEventAny::Canceled(_)))
);
if deferred {
for event in events.iter() {
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
}
let cache = cache.borrow();
assert_eq!(
cache.position(&position_id).unwrap().quantity,
Quantity::from("0.000")
);
for (id, quantity) in [("REENTRANT-A", "2.000"), ("REENTRANT-B", "8.000")] {
let order = cache.order(&ClientOrderId::from(id)).unwrap();
assert_eq!(order.status(), OrderStatus::Filled);
assert_eq!(order.quantity(), Quantity::from(quantity));
assert_eq!(order.filled_qty(), Quantity::from(quantity));
}
}
#[rstest]
fn test_position_sync_includes_newly_activated_oto_child(
#[values(false, true)] deferred: bool,
) {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let position_id = PositionId::from("ACTIVATION-POSITION");
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L2_MBP,
OmsType::Hedging,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
Default::default(),
);
let (opening, opening_fill) = pending_position_fill(
&instrument,
position_id,
"ACTIVATION-OPEN",
OrderSide::Buy,
"10.000",
);
let position = Position::new(&instrument, opening_fill.clone());
engine
.account_ids
.insert(position.trader_id, position.account_id);
cache
.borrow_mut()
.add_order(opening, Some(position_id), None, false)
.unwrap();
cache
.borrow_mut()
.update_order(&OrderEventAny::Filled(opening_fill))
.unwrap();
cache
.borrow_mut()
.add_position(&position, OmsType::Hedging)
.unwrap();
let parent_id = ClientOrderId::from("ACTIVATION-PARENT");
let child_id = ClientOrderId::from("ACTIVATION-CHILD");
let parent = OrderTestBuilder::new(OrderType::Market)
.instrument_id(instrument.id())
.client_order_id(parent_id)
.side(OrderSide::Buy)
.quantity(Quantity::from("10.000"))
.contingency_type(ContingencyType::Oto)
.linked_order_ids(vec![child_id])
.submit(true)
.build();
let child = OrderTestBuilder::new(OrderType::Limit)
.instrument_id(instrument.id())
.client_order_id(child_id)
.side(OrderSide::Sell)
.quantity(Quantity::from("10.000"))
.price(Price::from("2000.00"))
.reduce_only(true)
.parent_order_id(parent_id)
.submit(true)
.build();
for order in [parent.clone(), child] {
cache
.borrow_mut()
.add_order(order, Some(position_id), None, false)
.unwrap();
}
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = events.clone();
let handler_cache = cache.clone();
engine.set_event_handler(Rc::new(move |event| {
if !deferred || matches!(event, OrderEventAny::Accepted(_)) {
handler_cache.borrow_mut().update_order(&event).unwrap();
if let OrderEventAny::Filled(fill) = &event {
handler_cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
events_handler.borrow_mut().push(event);
}));
assert!(!engine.order_exists(child_id));
engine
.apply_fills(
&parent,
&[(Price::from("1000.00"), Quantity::from("2.000"))],
LiquiditySide::Taker,
Some(position_id),
Some(&position),
None,
)
.unwrap();
let events = events.borrow();
assert_eq!(events.len(), 3);
assert!(
matches!(&events[0], OrderEventAny::Filled(fill) if fill.client_order_id == parent_id && fill.last_qty == Quantity::from("2.000"))
);
assert!(
matches!(&events[1], OrderEventAny::Accepted(accepted) if accepted.client_order_id == child_id)
);
let OrderEventAny::Updated(update) = &events[2] else {
panic!("Expected child quantity update")
};
assert_eq!(update.client_order_id, child_id);
assert_eq!(update.quantity, Quantity::from("2.000"));
assert_eq!(update.price, Some(Price::from("2000.00")));
assert_eq!(update.trigger_price, None);
assert!(engine.order_exists(child_id));
assert_eq!(
engine.order_snapshot(child_id).unwrap().quantity(),
Quantity::from("2.000")
);
if deferred {
for event in events.iter() {
if matches!(event, OrderEventAny::Accepted(_)) {
continue;
}
cache.borrow_mut().update_order(event).unwrap();
if let OrderEventAny::Filled(fill) = event {
cache
.borrow_mut()
.update_position_from_fill(position_id, fill)
.unwrap();
}
}
}
let cache = cache.borrow();
assert_eq!(
cache.position(&position_id).unwrap().quantity,
Quantity::from("12.000")
);
let child = cache.order(&child_id).unwrap();
assert_eq!(child.status(), OrderStatus::Accepted);
assert_eq!(child.quantity(), Quantity::from("2.000"));
assert_eq!(child.filled_qty(), Quantity::from("0.000"));
assert_eq!(child.leaves_qty(), Quantity::from("2.000"));
}
fn pending_position_fill(
instrument: &InstrumentAny,
position_id: PositionId,
id: &str,
side: OrderSide,
quantity: &str,
) -> (OrderAny, OrderFilled) {
let order = OrderTestBuilder::new(OrderType::Market)
.instrument_id(instrument.id())
.client_order_id(ClientOrderId::from(id))
.side(side)
.quantity(Quantity::from(quantity))
.submit(true)
.build();
let OrderEventAny::Filled(fill) = TestOrderEventStubs::filled(
&order,
instrument,
Some(TradeId::from(id)),
Some(position_id),
Some(Price::from("1000.00")),
None,
None,
Some(Money::zero(instrument.quote_currency())),
None,
None,
) else {
unreachable!()
};
(order, fill)
}
#[rstest]
fn test_pending_modify_updates_acknowledge_individually_and_reset() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L2_MBP,
OmsType::Netting,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
cache.clone(),
OrderMatchingEngineConfig::default(),
);
let mut order = OrderTestBuilder::new(OrderType::Limit)
.instrument_id(instrument.id())
.side(OrderSide::Buy)
.quantity(Quantity::from("1.000"))
.price(Price::from("99.00"))
.submit(true)
.build();
let id = order.client_order_id();
engine.set_event_handler(Rc::new(|_| {}));
engine.process_order(&mut order, AccountId::from("ACCOUNT-001"));
let pending = Rc::new(RefCell::new(Vec::new()));
let events = pending.clone();
engine.set_event_handler(Rc::new(move |event| events.borrow_mut().push(event)));
for (quantity, price) in [
(Some(Quantity::from("2.000")), None),
(None, Some(Price::from("100.00"))),
] {
engine.process_modify(
&ModifyOrder::new(
order.trader_id(),
None,
order.strategy_id(),
order.instrument_id(),
id,
None,
quantity,
price,
None,
UUID4::new(),
UnixNanos::from(1),
None,
None,
),
AccountId::from("ACCOUNT-001"),
);
}
assert_eq!(pending.borrow().len(), 2);
cache
.borrow_mut()
.update_order(&pending.borrow()[0])
.unwrap();
let snapshot = engine.order_snapshot(id).unwrap();
assert_eq!(snapshot.quantity(), Quantity::from("2.000"));
assert_eq!(snapshot.price(), Some(Price::from("100.00")));
assert_eq!(engine.pending_order_updates.borrow()[&id].len(), 1);
cache
.borrow_mut()
.update_order(&pending.borrow()[1])
.unwrap();
engine.iterate(UnixNanos::from(2), AggressorSide::NoAggressor);
assert!(engine.pending_order_updates.borrow().is_empty());
engine.process_modify(
&ModifyOrder::new(
order.trader_id(),
None,
order.strategy_id(),
order.instrument_id(),
id,
None,
Some(Quantity::from("3.000")),
None,
None,
UUID4::new(),
UnixNanos::from(3),
None,
None,
),
AccountId::from("ACCOUNT-001"),
);
assert_eq!(
engine.order_snapshot(id).unwrap().quantity(),
Quantity::from("3.000")
);
engine.reset();
assert!(engine.pending_order_updates.borrow().is_empty());
assert_eq!(
engine.order_snapshot(id).unwrap().quantity(),
Quantity::from("2.000")
);
}
#[rstest]
fn test_process_order_rejects_reduce_only_when_support_is_disabled() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
Rc::new(RefCell::new(Cache::default())),
OrderMatchingEngineConfig::builder()
.use_reduce_only(false)
.build(),
);
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = Rc::clone(&events);
engine.set_event_handler(Rc::new(move |event| {
events_handler.borrow_mut().push(event);
}));
let mut order = OrderTestBuilder::new(OrderType::Market)
.instrument_id(instrument.id())
.side(OrderSide::Sell)
.quantity(Quantity::from("1.000"))
.reduce_only(true)
.submit(true)
.build();
engine.process_order(&mut order, AccountId::from("ACCOUNT-001"));
let events = events.borrow();
assert_eq!(events.len(), 1);
let OrderEventAny::Rejected(rejected) = &events[0] else {
panic!("Expected OrderRejected, was {:?}", events[0]);
};
assert_eq!(
rejected.reason,
"Reduce-only orders are not supported by this matching engine"
);
}
#[rstest]
fn test_post_match_order_action_does_not_clone_closed_order() {
let order = post_match_closed_limit_order();
let clone_count = Cell::new(0);
let action = post_match_order_action(&order, true, UnixNanos::from(1_u64), |order| {
clone_count.set(clone_count.get() + 1);
order.clone()
});
assert!(matches!(action, PostMatchOrderAction::RemoveClosed));
assert_eq!(clone_count.get(), 0);
}
#[rstest]
fn test_post_match_order_action_clones_expired_gtd_order_once() {
let order = post_match_gtd_limit_order();
let clone_count = Cell::new(0);
let action = post_match_order_action(&order, true, UnixNanos::from(10_u64), |order| {
clone_count.set(clone_count.get() + 1);
order.clone()
});
let PostMatchOrderAction::Expire(cloned) = action else {
panic!("Expected expired action, was {action:?}");
};
assert_eq!(cloned.client_order_id(), order.client_order_id());
assert_eq!(clone_count.get(), 1);
}
#[rstest]
fn test_post_match_order_action_clones_trailing_order_once() {
let order = post_match_trailing_stop_order();
let clone_count = Cell::new(0);
let action = post_match_order_action(&order, true, UnixNanos::from(1_u64), |order| {
clone_count.set(clone_count.get() + 1);
order.clone()
});
let PostMatchOrderAction::UpdateTrailing(cloned) = action else {
panic!("Expected trailing update action, was {action:?}");
};
assert_eq!(cloned.client_order_id(), order.client_order_id());
assert_eq!(clone_count.get(), 1);
}
fn post_match_limit_order() -> OrderAny {
OrderTestBuilder::new(OrderType::Limit)
.instrument_id(crypto_perpetual_ethusdt().id())
.side(OrderSide::Buy)
.price(Price::from("1500.00"))
.quantity(Quantity::from("1.000"))
.client_order_id(ClientOrderId::from("POST-MATCH-LIMIT"))
.submit(true)
.build()
}
fn post_match_closed_limit_order() -> OrderAny {
let account_id = AccountId::from("SIM-001");
let venue_order_id = VenueOrderId::from("V-001");
let mut order = post_match_limit_order();
order
.apply(TestOrderEventStubs::accepted(
&order,
account_id,
venue_order_id,
))
.unwrap();
order
.apply(TestOrderEventStubs::canceled(
&order,
account_id,
Some(venue_order_id),
))
.unwrap();
order
}
fn post_match_gtd_limit_order() -> OrderAny {
OrderTestBuilder::new(OrderType::Limit)
.instrument_id(crypto_perpetual_ethusdt().id())
.side(OrderSide::Buy)
.price(Price::from("1500.00"))
.quantity(Quantity::from("1.000"))
.time_in_force(TimeInForce::Gtd)
.expire_time(UnixNanos::from(10_u64))
.client_order_id(ClientOrderId::from("POST-MATCH-GTD"))
.submit(true)
.build()
}
fn post_match_trailing_stop_order() -> OrderAny {
OrderTestBuilder::new(OrderType::TrailingStopMarket)
.instrument_id(crypto_perpetual_ethusdt().id())
.side(OrderSide::Buy)
.quantity(Quantity::from("1.000"))
.trigger_price(Price::from("1510.00"))
.trigger_type(TriggerType::BidAsk)
.trailing_offset(Decimal::new(5, 0))
.trailing_offset_type(TrailingOffsetType::Price)
.client_order_id(ClientOrderId::from("POST-MATCH-TRAIL"))
.submit(true)
.build()
}
#[rstest]
fn test_fill_order_calculates_commission_from_fill_liquidity_side() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let clock = Rc::new(RefCell::new(TestClock::new()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
cache,
Default::default(),
);
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = Rc::clone(&events);
engine.set_event_handler(Rc::new(move |event| {
events_handler.borrow_mut().push(event);
}));
let mut order = OrderTestBuilder::new(OrderType::Market)
.instrument_id(instrument.id())
.side(OrderSide::Buy)
.quantity(Quantity::from("1.000"))
.submit(true)
.build();
order.set_liquidity_side(LiquiditySide::Maker);
engine
.account_ids
.insert(order.trader_id(), AccountId::from("ACCOUNT-001"));
engine
.fill_order(
&order,
Price::from("1500.00"),
Quantity::from("1.000"),
LiquiditySide::Taker,
None,
None,
)
.unwrap();
let events = events.borrow();
assert_eq!(events.len(), 1);
let fill = match &events[0] {
OrderEventAny::Filled(fill) => fill,
event => panic!("Expected OrderFilled, was {event:?}"),
};
let commission = fill.commission.expect("expected commission");
let expected_commission =
fill.last_qty.as_decimal() * fill.last_px.as_decimal() * instrument.taker_fee();
assert_eq!(fill.liquidity_side, LiquiditySide::Taker);
assert_eq!(commission.currency, instrument.quote_currency());
assert_eq!(commission.as_decimal(), expected_commission);
}
#[rstest]
fn test_custom_fee_model_handle_is_called_by_fill_order() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let clock = Rc::new(RefCell::new(TestClock::new()));
let calls = Rc::new(Cell::new(0));
let expected_commission = Money::from("1.23 USDT");
let fee_model = FeeModelHandle::new(RecordingFeeModel {
calls: Rc::clone(&calls),
commission: expected_commission,
});
let cloned_fee_model = fee_model.clone();
drop(fee_model);
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
cloned_fee_model,
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
cache,
Default::default(),
);
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = Rc::clone(&events);
engine.set_event_handler(Rc::new(move |event| {
events_handler.borrow_mut().push(event);
}));
let order = OrderTestBuilder::new(OrderType::Market)
.instrument_id(instrument.id())
.side(OrderSide::Buy)
.quantity(Quantity::from("1.000"))
.submit(true)
.build();
engine
.account_ids
.insert(order.trader_id(), AccountId::from("ACCOUNT-001"));
engine
.fill_order(
&order,
Price::from("1500.00"),
Quantity::from("1.000"),
LiquiditySide::Taker,
None,
None,
)
.unwrap();
let events = events.borrow();
assert_eq!(events.len(), 1);
let fill = match &events[0] {
OrderEventAny::Filled(fill) => fill,
event => panic!("Expected OrderFilled, was {event:?}"),
};
assert_eq!(calls.get(), 1);
assert_eq!(fill.commission, Some(expected_commission));
}
#[rstest]
fn test_fill_order_does_not_cache_filled_qty_when_fee_model_fails() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let clock = Rc::new(RefCell::new(TestClock::new()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelHandle::new(FailingFeeModel),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
cache,
Default::default(),
);
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = Rc::clone(&events);
engine.set_event_handler(Rc::new(move |event| {
events_handler.borrow_mut().push(event);
}));
let order = OrderTestBuilder::new(OrderType::Market)
.instrument_id(instrument.id())
.side(OrderSide::Buy)
.quantity(Quantity::from("1.000"))
.submit(true)
.build();
engine
.account_ids
.insert(order.trader_id(), AccountId::from("ACCOUNT-001"));
let result = engine.fill_order(
&order,
Price::from("1500.00"),
Quantity::from("1.000"),
LiquiditySide::Taker,
None,
None,
);
assert!(result.is_err());
assert_eq!(engine.cached_filled_qty_len(), 0);
assert!(events.borrow().is_empty());
}
#[rstest]
fn test_process_cancel_all_includes_submitted_orders_for_selected_account() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let instrument_id = instrument.id();
let cache = Rc::new(RefCell::new(Cache::default()));
let clock = Rc::new(RefCell::new(TestClock::new()));
let mut engine = OrderMatchingEngine::new(
instrument,
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
Rc::clone(&cache),
Default::default(),
);
let selected_account = AccountId::from("ACCOUNT-001");
let other_account = AccountId::from("ACCOUNT-002");
let selected_strategy = StrategyId::from("STRATEGY-001");
let other_strategy = StrategyId::from("STRATEGY-002");
let selected_order = OrderTestBuilder::new(OrderType::Limit)
.strategy_id(selected_strategy)
.instrument_id(instrument_id)
.client_order_id(ClientOrderId::from("O-SUBMITTED-SELECTED"))
.side(OrderSide::Buy)
.price(Price::from("1400.00"))
.quantity(Quantity::from("1.000"))
.build();
let other_order = OrderTestBuilder::new(OrderType::Limit)
.strategy_id(other_strategy)
.instrument_id(instrument_id)
.client_order_id(ClientOrderId::from("O-SUBMITTED-OTHER"))
.side(OrderSide::Buy)
.price(Price::from("1300.00"))
.quantity(Quantity::from("1.000"))
.build();
{
let mut cache = cache.borrow_mut();
cache
.add_order(selected_order.clone(), None, None, false)
.unwrap();
cache
.add_order(other_order.clone(), None, None, false)
.unwrap();
cache
.update_order(&TestOrderEventStubs::submitted(
&selected_order,
selected_account,
))
.unwrap();
cache
.update_order(&TestOrderEventStubs::submitted(&other_order, other_account))
.unwrap();
}
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = Rc::clone(&events);
let event_cache = Rc::clone(&cache);
engine.set_event_handler(Rc::new(move |event| {
event_cache.borrow_mut().update_order(&event).unwrap();
events_handler.borrow_mut().push(event);
}));
let command = CancelAllOrders::new(
TraderId::from("TRADER-001"),
None,
StrategyId::from("CALLER-001"),
instrument_id,
None,
UUID4::new(),
UnixNanos::default(),
None,
None,
);
engine.process_cancel_all(&command, selected_account);
let events = events.borrow();
assert_eq!(events.len(), 1);
let OrderEventAny::Canceled(canceled) = &events[0] else {
panic!("Expected OrderCanceled, was {:?}", events[0]);
};
assert_eq!(canceled.client_order_id, selected_order.client_order_id());
assert_eq!(canceled.strategy_id, selected_strategy);
assert_eq!(canceled.account_id, Some(selected_account));
let cache = cache.borrow();
assert_eq!(
cache
.order(&selected_order.client_order_id())
.unwrap()
.status(),
OrderStatus::Canceled
);
assert_eq!(
cache
.order(&other_order.client_order_id())
.unwrap()
.status(),
OrderStatus::Submitted
);
}
#[rstest]
fn test_process_cancel_all_excluding_leaves_excluded_orders_untouched() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let instrument_id = instrument.id();
let cache = Rc::new(RefCell::new(Cache::default()));
let clock = Rc::new(RefCell::new(TestClock::new()));
let mut engine = OrderMatchingEngine::new(
instrument,
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
Rc::clone(&cache),
Default::default(),
);
let account_id = AccountId::from("ACCOUNT-001");
let strategy_id = StrategyId::from("STRATEGY-001");
let received = OrderTestBuilder::new(OrderType::Limit)
.strategy_id(strategy_id)
.instrument_id(instrument_id)
.client_order_id(ClientOrderId::from("O-RECEIVED"))
.side(OrderSide::Buy)
.price(Price::from("1400.00"))
.quantity(Quantity::from("1.000"))
.build();
let in_transit = OrderTestBuilder::new(OrderType::Limit)
.strategy_id(strategy_id)
.instrument_id(instrument_id)
.client_order_id(ClientOrderId::from("O-IN-TRANSIT"))
.side(OrderSide::Buy)
.price(Price::from("1300.00"))
.quantity(Quantity::from("1.000"))
.build();
{
let mut cache = cache.borrow_mut();
cache
.add_order(received.clone(), None, None, false)
.unwrap();
cache
.add_order(in_transit.clone(), None, None, false)
.unwrap();
cache
.update_order(&TestOrderEventStubs::submitted(&received, account_id))
.unwrap();
cache
.update_order(&TestOrderEventStubs::submitted(&in_transit, account_id))
.unwrap();
}
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = Rc::clone(&events);
let event_cache = Rc::clone(&cache);
engine.set_event_handler(Rc::new(move |event| {
event_cache.borrow_mut().update_order(&event).unwrap();
events_handler.borrow_mut().push(event);
}));
let command = CancelAllOrders::new(
TraderId::from("TRADER-001"),
None,
StrategyId::from("CALLER-001"),
instrument_id,
None,
UUID4::new(),
UnixNanos::default(),
None,
None,
);
engine.process_cancel_all_excluding(&command, account_id, &[in_transit.client_order_id()]);
let events = events.borrow();
assert_eq!(
events.len(),
1,
"expected one OrderCanceled, was {events:?}"
);
let OrderEventAny::Canceled(canceled) = &events[0] else {
panic!("Expected OrderCanceled, was {:?}", events[0]);
};
assert_eq!(canceled.client_order_id, received.client_order_id());
assert_eq!(
cache
.borrow()
.order(&in_transit.client_order_id())
.unwrap()
.status(),
OrderStatus::Submitted,
"an excluded order must be left untouched",
);
}
#[rstest]
fn test_process_cancel_all_excluding_spares_an_excluded_contingent_order() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let instrument_id = instrument.id();
let cache = Rc::new(RefCell::new(Cache::default()));
let clock = Rc::new(RefCell::new(TestClock::new()));
let mut engine = OrderMatchingEngine::new(
instrument,
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
Rc::clone(&cache),
Default::default(),
);
assert!(engine.config.support_contingent_orders);
let account_id = AccountId::from("ACCOUNT-001");
let strategy_id = StrategyId::from("STRATEGY-001");
let received_id = ClientOrderId::from("O-RECEIVED");
let in_transit_id = ClientOrderId::from("O-IN-TRANSIT");
let received = OrderTestBuilder::new(OrderType::Limit)
.strategy_id(strategy_id)
.instrument_id(instrument_id)
.client_order_id(received_id)
.side(OrderSide::Buy)
.price(Price::from("1400.00"))
.quantity(Quantity::from("1.000"))
.contingency_type(ContingencyType::Oco)
.linked_order_ids(vec![in_transit_id])
.build();
let in_transit = OrderTestBuilder::new(OrderType::Limit)
.strategy_id(strategy_id)
.instrument_id(instrument_id)
.client_order_id(in_transit_id)
.side(OrderSide::Buy)
.price(Price::from("1300.00"))
.quantity(Quantity::from("1.000"))
.contingency_type(ContingencyType::Oco)
.linked_order_ids(vec![received_id])
.build();
{
let mut cache = cache.borrow_mut();
cache
.add_order(received.clone(), None, None, false)
.unwrap();
cache
.add_order(in_transit.clone(), None, None, false)
.unwrap();
cache
.update_order(&TestOrderEventStubs::submitted(&received, account_id))
.unwrap();
cache
.update_order(&TestOrderEventStubs::submitted(&in_transit, account_id))
.unwrap();
}
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = Rc::clone(&events);
let event_cache = Rc::clone(&cache);
engine.set_event_handler(Rc::new(move |event| {
event_cache.borrow_mut().update_order(&event).unwrap();
events_handler.borrow_mut().push(event);
}));
let command = CancelAllOrders::new(
TraderId::from("TRADER-001"),
None,
StrategyId::from("CALLER-001"),
instrument_id,
None,
UUID4::new(),
UnixNanos::default(),
None,
None,
);
engine.process_cancel_all_excluding(&command, account_id, &[in_transit_id]);
let events = events.borrow();
assert_eq!(
events.len(),
1,
"expected one OrderCanceled, was {events:?}"
);
let OrderEventAny::Canceled(canceled) = &events[0] else {
panic!("Expected OrderCanceled, was {:?}", events[0]);
};
assert_eq!(canceled.client_order_id, received_id);
assert_eq!(
cache.borrow().order(&in_transit_id).unwrap().status(),
OrderStatus::Submitted,
"canceling its OCO sibling must not cancel an excluded order",
);
}
fn collision_engine() -> (OrderMatchingEngine, Rc<RefCell<Cache>>, VenueOrderId) {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let venue_order_id = VenueOrderId::from(format!("{}-1-1", instrument.id().venue));
cache
.borrow_mut()
.add_venue_order_id(&ClientOrderId::from("O-OWNER"), &venue_order_id, false)
.unwrap();
let engine = OrderMatchingEngine::new(
instrument,
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
Rc::new(RefCell::new(TestClock::new())),
Rc::clone(&cache),
Default::default(),
);
(engine, cache, venue_order_id)
}
#[rstest]
#[case(OrderType::Market)]
#[case(OrderType::MarketToLimit)]
fn test_market_collision_probes_and_fills_with_default_ack_config(
#[case] order_type: OrderType,
) {
let (mut engine, cache, venue_order_id) = collision_engine();
assert!(!engine.config.use_market_order_acks);
let quote = QuoteTick::new(
engine.instrument.id(),
Price::from("1499.00"),
Price::from("1500.00"),
Quantity::from("10.000"),
Quantity::from("10.000"),
UnixNanos::default(),
UnixNanos::default(),
);
engine.process_quote_tick("e);
let events = Rc::new(RefCell::new(Vec::new()));
let events_handler = Rc::clone(&events);
engine.set_event_handler(Rc::new(move |event| {
events_handler.borrow_mut().push(event);
}));
let mut order = OrderTestBuilder::new(order_type)
.instrument_id(engine.instrument.id())
.client_order_id(ClientOrderId::from("O-CLAIMANT"))
.side(OrderSide::Buy)
.quantity(Quantity::from("1.000"))
.submit(true)
.build();
engine.process_order(&mut order, AccountId::from("ACCOUNT-001"));
assert!(
!events
.borrow()
.iter()
.any(|event| matches!(event, OrderEventAny::Rejected(_)))
);
assert!(
events
.borrow()
.iter()
.any(|event| matches!(event, OrderEventAny::Filled(_)))
);
assert!(cache.borrow().order_exists(&order.client_order_id()));
assert_eq!(
cache.borrow().client_order_id(&venue_order_id),
Some(&ClientOrderId::from("O-OWNER"))
);
assert_eq!(
cache.borrow().venue_order_id(&order.client_order_id()),
Some(&VenueOrderId::from(format!("{}-1-2", engine.venue)))
);
}
struct RecordingFeeModel {
calls: Rc<Cell<u32>>,
commission: Money,
}
impl FeeModel for RecordingFeeModel {
fn get_commission(
&self,
_order: &OrderAny,
_fill_quantity: Quantity,
_fill_px: Price,
_instrument: &InstrumentAny,
) -> anyhow::Result<Money> {
self.calls.set(self.calls.get() + 1);
Ok(self.commission)
}
}
struct FailingFeeModel;
impl FeeModel for FailingFeeModel {
fn get_commission(
&self,
_order: &OrderAny,
_fill_quantity: Quantity,
_fill_px: Price,
_instrument: &InstrumentAny,
) -> anyhow::Result<Money> {
Err(anyhow::anyhow!("fee model failed"))
}
}
#[rstest]
fn test_custom_fill_model_handle_is_called_by_market_fill() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let clock = Rc::new(RefCell::new(TestClock::new()));
let calls = Rc::new(Cell::new(0));
let fill_model = FillModelHandle::new(RecordingFillModel {
calls: Rc::clone(&calls),
});
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
fill_model,
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
cache,
Default::default(),
);
let quote = QuoteTick::new(
instrument.id(),
Price::from("1500.00"),
Price::from("1501.00"),
Quantity::from("10.000"),
Quantity::from("10.000"),
UnixNanos::default(),
UnixNanos::default(),
);
engine.process_quote_tick("e);
let mut order = OrderTestBuilder::new(OrderType::Market)
.instrument_id(instrument.id())
.side(OrderSide::Buy)
.quantity(Quantity::from("1.000"))
.submit(true)
.build();
engine.process_order(&mut order, AccountId::from("ACCOUNT-001"));
assert_eq!(calls.get(), 1);
}
#[rstest]
fn test_l1_depth10_skips_padding_for_last_quote_tracking() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let cache = Rc::new(RefCell::new(Cache::default()));
let clock = Rc::new(RefCell::new(TestClock::new()));
let mut engine = OrderMatchingEngine::new(
instrument.clone(),
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
cache,
Default::default(),
);
let mut bids = [BookOrder::default(); DEPTH10_LEN];
let mut asks = [BookOrder::default(); DEPTH10_LEN];
bids[1] = BookOrder::new(
OrderSide::Buy,
Price::from("1499.00"),
Quantity::from("1.000"),
1,
);
asks[0] = BookOrder::new(
OrderSide::Sell,
Price::from("1500.00"),
Quantity::from("1.000"),
2,
);
let depth = OrderBookDepth10::new(
instrument.id(),
bids,
asks,
[0; DEPTH10_LEN],
[0; DEPTH10_LEN],
0,
0,
UnixNanos::from(1_u64),
UnixNanos::from(1_u64),
);
engine.process_order_book_depth10(&depth).unwrap();
assert_eq!(engine.last_quote_bid, Some(Price::from("1499.00")));
assert_eq!(engine.last_quote_ask, Some(Price::from("1500.00")));
let depth_without_bid = OrderBookDepth10::new(
instrument.id(),
[BookOrder::default(); DEPTH10_LEN],
asks,
[0; DEPTH10_LEN],
[0; DEPTH10_LEN],
0,
1,
UnixNanos::from(2_u64),
UnixNanos::from(2_u64),
);
engine
.process_order_book_depth10(&depth_without_bid)
.unwrap();
assert_eq!(engine.last_quote_bid, None);
assert_eq!(engine.last_quote_ask, Some(Price::from("1500.00")));
}
struct RecordingFillModel {
calls: Rc<Cell<u32>>,
}
impl FillModel for RecordingFillModel {
fn is_limit_filled(&mut self) -> anyhow::Result<bool> {
Ok(true)
}
fn is_slipped(&mut self) -> anyhow::Result<bool> {
Ok(false)
}
fn get_orderbook_for_fill_simulation(
&mut self,
_instrument: &InstrumentAny,
_order: &OrderAny,
_best_bid: Price,
_best_ask: Price,
) -> anyhow::Result<Option<OrderBook>> {
self.calls.set(self.calls.get() + 1);
Ok(None)
}
}
#[rstest]
fn test_fee_underlying_price_uses_valid_cached_greeks_price() {
let instrument = InstrumentAny::CryptoOption(crypto_option_btc_deribit(
3,
1,
Price::from("0.001"),
Quantity::from("0.1"),
));
let cache = Rc::new(RefCell::new(Cache::default()));
cache.borrow_mut().add_option_greeks(OptionGreeks {
instrument_id: instrument.id(),
underlying_price: Some(50_000.0),
..Default::default()
});
let clock = Rc::new(RefCell::new(TestClock::new()));
let engine = OrderMatchingEngine::new(
instrument,
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
cache,
Default::default(),
);
let price = engine
.fee_underlying_price()
.unwrap()
.expect("expected underlying price");
assert_eq!(price.precision, FIXED_PRECISION);
assert_eq!(price.as_decimal(), Decimal::from(50_000));
}
#[rstest]
fn test_fee_underlying_price_rejects_invalid_cached_greeks_price() {
let instrument = InstrumentAny::CryptoOption(crypto_option_btc_deribit(
3,
1,
Price::from("0.001"),
Quantity::from("0.1"),
));
let cache = Rc::new(RefCell::new(Cache::default()));
cache.borrow_mut().add_option_greeks(OptionGreeks {
instrument_id: instrument.id(),
underlying_price: Some(f64::NAN),
..Default::default()
});
let clock = Rc::new(RefCell::new(TestClock::new()));
let engine = OrderMatchingEngine::new(
instrument,
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
BookType::L1_MBP,
OmsType::Netting,
AccountType::Margin,
clock,
cache,
Default::default(),
);
let error = engine.fee_underlying_price().unwrap_err();
assert_eq!(
error,
CorrectnessError::InvalidValue {
param: "value".to_string(),
value: "NaN".to_string(),
type_name: "f64",
}
);
}
#[rstest]
fn test_bar_tick_sizes_divisible() {
let volume = Quantity::from("100.000");
let increment = Quantity::from("0.001");
let sizes = BarTickSizes::from_volume(volume, increment);
assert_eq!(sizes.open, Quantity::from("25.000"));
assert_eq!(sizes.high, Quantity::from("25.000"));
assert_eq!(sizes.low, Quantity::from("25.000"));
assert_eq!(sizes.close, Quantity::from("25.000"));
assert_valid_bar_tick_sizes(volume, increment);
}
#[rstest]
fn test_bar_tick_sizes_indivisible_with_remainder() {
let volume = Quantity::from("0.05");
let increment = Quantity::from("0.01");
let sizes = BarTickSizes::from_volume(volume, increment);
assert_eq!(sizes.open, Quantity::from("0.01"));
assert_eq!(sizes.high, Quantity::from("0.01"));
assert_eq!(sizes.low, Quantity::from("0.01"));
assert_eq!(sizes.close, Quantity::from("0.02"));
assert_valid_bar_tick_sizes(volume, increment);
assert_eq!(
sizes.open.raw() + sizes.high.raw() + sizes.low.raw() + sizes.close.raw(),
volume.raw()
);
}
#[rstest]
#[case("1", "0", "0", "0", "1")]
#[case("2", "0", "1", "1", "0")]
#[case("3", "1", "1", "1", "0")]
fn test_bar_tick_sizes_units_less_than_four_preserves_volume(
#[case] volume: &str,
#[case] open_size: &str,
#[case] high_size: &str,
#[case] low_size: &str,
#[case] close_size: &str,
) {
let volume = Quantity::from(volume);
let increment = Quantity::from("1");
let sizes = BarTickSizes::from_volume(volume, increment);
assert_eq!(sizes.open, Quantity::from(open_size));
assert_eq!(sizes.high, Quantity::from(high_size));
assert_eq!(sizes.low, Quantity::from(low_size));
assert_eq!(sizes.close, Quantity::from(close_size));
assert_valid_bar_tick_sizes(volume, increment);
assert_eq!(
sizes.open.raw() + sizes.high.raw() + sizes.low.raw() + sizes.close.raw(),
volume.raw()
);
}
#[rstest]
fn test_bar_tick_sizes_zero_volume_remains_zero() {
let volume = Quantity::zero(3);
let increment = Quantity::from("0.001");
let sizes = BarTickSizes::from_volume(volume, increment);
assert_eq!(sizes.open, Quantity::zero(3));
assert_eq!(sizes.high, Quantity::zero(3));
assert_eq!(sizes.low, Quantity::zero(3));
assert_eq!(sizes.close, Quantity::zero(3));
assert_valid_bar_tick_sizes(volume, increment);
}
#[rstest]
fn test_bar_tick_sizes_rounds_down_to_size_increment() {
let volume = Quantity::from("1.07");
let increment = Quantity::from("0.10");
let sizes = BarTickSizes::from_volume(volume, increment);
assert_eq!(sizes.open, Quantity::from("0.20"));
assert_eq!(sizes.high, Quantity::from("0.20"));
assert_eq!(sizes.low, Quantity::from("0.20"));
assert_eq!(sizes.close, Quantity::from("0.40"));
assert_valid_bar_tick_sizes(volume, increment);
}
#[rstest]
fn test_bar_tick_sizes_at_fixed_precision() {
let units: QuantityRaw = 17;
let volume = Quantity::from_raw(units, FIXED_PRECISION);
let increment = Quantity::from_raw(1, FIXED_PRECISION);
let sizes = BarTickSizes::from_volume(volume, increment);
assert_eq!(sizes.open.raw(), 4);
assert_eq!(sizes.high.raw(), 4);
assert_eq!(sizes.low.raw(), 4);
assert_eq!(sizes.close.raw(), 5);
assert_valid_bar_tick_sizes(volume, increment);
}
fn get_queue_engine(
instrument: InstrumentAny,
book_type: BookType,
) -> (OrderMatchingEngine, Rc<RefCell<Cache>>) {
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let config = OrderMatchingEngineConfig {
trade_execution: true,
queue_position: true,
..Default::default()
};
let mut engine = OrderMatchingEngine::new(
instrument,
1,
FillModelHandle::default(),
FeeModelAny::default().into(),
book_type,
OmsType::Netting,
AccountType::Margin,
clock,
Rc::clone(&cache),
config,
);
let handler_cache = Rc::clone(&cache);
engine.set_event_handler(Rc::new(move |event: OrderEventAny| {
if let Ok(mut cache) = handler_cache.try_borrow_mut() {
let _ = cache.update_order(&event);
}
}));
(engine, cache)
}
fn get_l3_queue_engine(instrument: InstrumentAny) -> (OrderMatchingEngine, Rc<RefCell<Cache>>) {
get_queue_engine(instrument, BookType::L3_MBO)
}
fn assert_l3_queue_synced(engine: &OrderMatchingEngine) {
for (client_order_id, orders_ahead) in &engine.queue_ahead_orders {
let set_sum: QuantityRaw = orders_ahead.values().sum();
let counter = engine
.queue_ahead_total
.get(client_order_id)
.map_or(0, |&(_, ahead_raw)| ahead_raw);
assert_eq!(
set_sum, counter,
"tracked orders out of sync with quantity-ahead counter for {client_order_id}",
);
}
for (client_order_id, price_raw) in &engine.queue_pending {
assert!(
engine
.queue_ids_by_price
.get(price_raw)
.is_some_and(|ids| ids.contains(client_order_id)),
"pending order {client_order_id} missing from price index",
);
}
for (client_order_id, (price_raw, _)) in &engine.queue_ahead_total {
assert!(
engine
.queue_ids_by_price
.get(price_raw)
.is_some_and(|ids| ids.contains(client_order_id)),
"tracked order {client_order_id} missing from price index",
);
}
for (price_raw, client_order_ids) in &engine.queue_ids_by_price {
for client_order_id in client_order_ids {
let pending_at_price = engine.queue_pending.get(client_order_id) == Some(price_raw);
let tracked_at_price = engine
.queue_ahead_total
.get(client_order_id)
.is_some_and(|(tracked_price_raw, _)| tracked_price_raw == price_raw);
assert!(
pending_at_price || tracked_at_price,
"price index contains stale order {client_order_id}",
);
}
}
}
#[rstest]
fn test_reset_clears_queue_positions() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let (mut engine, _cache) = get_l3_queue_engine(instrument);
let price = Price::from("100.00");
let client_order_id = ClientOrderId::from("O-RESET-QUEUE");
rest_l3_queue_order(&mut engine, price, 1, client_order_id);
assert!(engine.queue_ahead_total.contains_key(&client_order_id));
assert!(engine.queue_ahead_orders.contains_key(&client_order_id));
assert!(
engine
.queue_ids_by_price
.get(&price.raw())
.is_some_and(|ids| ids.contains(&client_order_id)),
);
engine.reset();
assert!(engine.queue_pending.is_empty());
assert!(engine.queue_ahead_total.is_empty());
assert!(engine.queue_ahead_orders.is_empty());
assert!(engine.queue_excess.is_empty());
assert!(engine.queue_ids_by_price.is_empty());
}
#[rstest]
fn test_cancel_removes_queue_position() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let (mut engine, _cache) = get_l3_queue_engine(instrument);
let price = Price::from("100.00");
let order =
rest_l3_queue_order(&mut engine, price, 1, ClientOrderId::from("O-CANCEL-QUEUE"));
let client_order_id = order.client_order_id();
assert!(engine.queue_ahead_total.contains_key(&client_order_id));
assert!(engine.queue_ahead_orders.contains_key(&client_order_id));
assert!(
engine
.queue_ids_by_price
.get(&price.raw())
.is_some_and(|ids| ids.contains(&client_order_id)),
);
engine.cancel_order(&order, None);
assert!(!engine.queue_pending.contains_key(&client_order_id));
assert!(!engine.queue_ahead_total.contains_key(&client_order_id));
assert!(!engine.queue_ahead_orders.contains_key(&client_order_id));
assert!(!engine.queue_excess.contains_key(&client_order_id));
assert!(!engine.queue_ids_by_price.contains_key(&price.raw()));
}
#[rstest]
fn test_modify_reindexes_queue_position() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let (mut engine, _cache) = get_l3_queue_engine(instrument);
let old_price = Price::from("100.00");
let new_price = Price::from("101.00");
let client_order_id = ClientOrderId::from("O-MODIFY-QUEUE");
let order = rest_l3_queue_order(&mut engine, old_price, 1, client_order_id);
let new_level = OrderBookDelta::new(
engine.instrument.id(),
BookAction::Add,
BookOrder::new(OrderSide::Sell, new_price, Quantity::from("10.000"), 2),
0,
2,
UnixNanos::from(2),
UnixNanos::from(2),
);
engine.process_order_book_delta(&new_level).unwrap();
let command = ModifyOrder::new(
order.trader_id(),
None,
order.strategy_id(),
order.instrument_id(),
client_order_id,
order.venue_order_id(),
None,
Some(new_price),
None,
UUID4::new(),
UnixNanos::from(3),
None,
None,
);
engine.process_modify(&command, AccountId::from("SIM-001"));
assert!(!engine.queue_ids_by_price.contains_key(&old_price.raw()));
assert_eq!(
engine
.queue_ids_by_price
.get(&new_price.raw())
.map(|ids| ids.iter().copied().collect::<Vec<_>>()),
Some(vec![client_order_id]),
);
assert_eq!(
engine.queue_ahead_total.get(&client_order_id),
Some(&(new_price.raw(), Quantity::from("10.000").raw())),
);
assert_eq!(
engine
.queue_ahead_orders
.get(&client_order_id)
.map(|orders| orders.keys().copied().collect::<Vec<_>>()),
Some(vec![2]),
);
}
#[rstest]
fn test_snapshot_rebases_l2_queue_position_after_size_decrease() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let instrument_id = instrument.id();
let (mut engine, cache) = get_queue_engine(instrument, BookType::L2_MBP);
let initial = OrderBookDelta::new(
instrument_id,
BookAction::Add,
BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("10.000"),
0,
),
0,
1,
UnixNanos::from(1_u64),
UnixNanos::from(1_u64),
);
engine.process_order_book_delta(&initial).unwrap();
let client_order_id = ClientOrderId::from("O-SNAPSHOT-DECREASE");
let mut order = OrderTestBuilder::new(OrderType::Limit)
.instrument_id(instrument_id)
.side(OrderSide::Sell)
.price(Price::from("100.00"))
.quantity(Quantity::from("1.000"))
.client_order_id(client_order_id)
.submit(true)
.build();
engine.process_order(&mut order, AccountId::from("SIM-001"));
assert_eq!(
engine.queue_ahead_total.get(&client_order_id),
Some(&(Price::from("100.00").raw(), Quantity::from("10.000").raw())),
);
let clear = OrderBookDelta::clear(
instrument_id,
2,
UnixNanos::from(2_u64),
UnixNanos::from(2_u64),
);
engine.process_order_book_delta(&clear).unwrap();
assert_eq!(
engine.queue_ahead_total.get(&client_order_id),
Some(&(Price::from("100.00").raw(), Quantity::from("10.000").raw())),
"partial snapshot must not discard the old queue estimate",
);
let snapshot = OrderBookDelta::new(
instrument_id,
BookAction::Add,
BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("8.000"),
0,
),
RecordFlag::F_LAST as u8,
2,
UnixNanos::from(2_u64),
UnixNanos::from(2_u64),
);
engine.process_order_book_delta(&snapshot).unwrap();
assert_eq!(
engine.queue_ahead_total.get(&client_order_id),
Some(&(Price::from("100.00").raw(), Quantity::from("8.000").raw())),
);
assert!(cache.borrow().order(&client_order_id).is_some());
}
#[rstest]
fn test_snapshot_rebase_does_not_increase_l2_queue_position() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let instrument_id = instrument.id();
let (mut engine, _cache) = get_queue_engine(instrument, BookType::L2_MBP);
let initial = OrderBookDelta::new(
instrument_id,
BookAction::Add,
BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("10.000"),
0,
),
0,
1,
UnixNanos::from(1_u64),
UnixNanos::from(1_u64),
);
engine.process_order_book_delta(&initial).unwrap();
let client_order_id = ClientOrderId::from("O-SNAPSHOT-INCREASE");
let mut order = OrderTestBuilder::new(OrderType::Limit)
.instrument_id(instrument_id)
.side(OrderSide::Sell)
.price(Price::from("100.00"))
.quantity(Quantity::from("1.000"))
.client_order_id(client_order_id)
.submit(true)
.build();
engine.process_order(&mut order, AccountId::from("SIM-001"));
let snapshot = OrderBookDelta::new(
instrument_id,
BookAction::Add,
BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("15.000"),
0,
),
RecordFlag::F_SNAPSHOT as u8 | RecordFlag::F_LAST as u8,
2,
UnixNanos::from(2_u64),
UnixNanos::from(2_u64),
);
engine.process_order_book_delta(&snapshot).unwrap();
assert_eq!(
engine.queue_ahead_total.get(&client_order_id),
Some(&(Price::from("100.00").raw(), Quantity::from("10.000").raw())),
);
}
#[rstest]
fn test_depth10_rebases_l2_queue_position() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let instrument_id = instrument.id();
let (mut engine, _cache) = get_queue_engine(instrument, BookType::L2_MBP);
let mut asks = [BookOrder::default(); DEPTH10_LEN];
asks[0] = BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("10.000"),
0,
);
let initial = OrderBookDepth10::new(
instrument_id,
[BookOrder::default(); DEPTH10_LEN],
asks,
[0; DEPTH10_LEN],
[0; DEPTH10_LEN],
0,
1,
UnixNanos::from(1_u64),
UnixNanos::from(1_u64),
);
engine.process_order_book_depth10(&initial).unwrap();
let client_order_id = ClientOrderId::from("O-DEPTH10-REBASE");
let mut order = OrderTestBuilder::new(OrderType::Limit)
.instrument_id(instrument_id)
.side(OrderSide::Sell)
.price(Price::from("100.00"))
.quantity(Quantity::from("1.000"))
.client_order_id(client_order_id)
.submit(true)
.build();
engine.process_order(&mut order, AccountId::from("SIM-001"));
asks[0] = BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("8.000"),
0,
);
let replacement = OrderBookDepth10::new(
instrument_id,
[BookOrder::default(); DEPTH10_LEN],
asks,
[0; DEPTH10_LEN],
[0; DEPTH10_LEN],
0,
2,
UnixNanos::from(2_u64),
UnixNanos::from(2_u64),
);
engine.process_order_book_depth10(&replacement).unwrap();
assert_eq!(
engine.queue_ahead_total.get(&client_order_id),
Some(&(Price::from("100.00").raw(), Quantity::from("8.000").raw())),
);
}
#[rstest]
fn test_snapshot_rebases_each_l3_order_independently() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let instrument_id = instrument.id();
let (mut engine, _cache) = get_l3_queue_engine(instrument);
for (order_id, sequence) in [(1, 1), (2, 2)] {
let delta = OrderBookDelta::new(
instrument_id,
BookAction::Add,
BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("5.000"),
order_id,
),
0,
sequence,
UnixNanos::from(sequence),
UnixNanos::from(sequence),
);
engine.process_order_book_delta(&delta).unwrap();
}
let client_order_id = ClientOrderId::from("O-SNAPSHOT-L3");
let order = rest_l3_queue_order(&mut engine, Price::from("100.00"), 3, client_order_id);
assert_eq!(
engine.queue_ahead_orders[&client_order_id]
.keys()
.copied()
.collect::<Vec<_>>(),
vec![1, 2, 3],
);
assert_eq!(
engine.queue_ahead_total.get(&client_order_id),
Some(&(Price::from("100.00").raw(), Quantity::from("20.000").raw())),
);
let snapshot = OrderBookDeltas::new(
instrument_id,
vec![
OrderBookDelta::clear(
instrument_id,
4,
UnixNanos::from(4_u64),
UnixNanos::from(4_u64),
),
OrderBookDelta::new(
instrument_id,
BookAction::Add,
BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("10.000"),
1,
),
RecordFlag::F_SNAPSHOT as u8,
4,
UnixNanos::from(4_u64),
UnixNanos::from(4_u64),
),
OrderBookDelta::new(
instrument_id,
BookAction::Add,
BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("5.000"),
2,
),
RecordFlag::F_LAST as u8,
4,
UnixNanos::from(4_u64),
UnixNanos::from(4_u64),
),
],
);
engine.process_order_book_deltas(&snapshot).unwrap();
assert_eq!(
engine.queue_ahead_orders[&client_order_id]
.iter()
.map(|(&order_id, &size)| (order_id, size))
.collect::<Vec<_>>(),
vec![
(1, Quantity::from("5.000").raw()),
(2, Quantity::from("5.000").raw())
],
);
assert_eq!(
engine.queue_ahead_total.get(&client_order_id),
Some(&(Price::from("100.00").raw(), Quantity::from("10.000").raw())),
);
let delete_a = OrderBookDelta::new(
instrument_id,
BookAction::Delete,
BookOrder::new(
OrderSide::Sell,
Price::from("100.00"),
Quantity::from("10.000"),
1,
),
0,
5,
UnixNanos::from(5_u64),
UnixNanos::from(5_u64),
);
engine.process_order_book_delta(&delete_a).unwrap();
assert_eq!(
engine.queue_ahead_orders[&client_order_id]
.keys()
.copied()
.collect::<Vec<_>>(),
vec![2],
);
assert_eq!(
engine.queue_ahead_total.get(&client_order_id),
Some(&(Price::from("100.00").raw(), Quantity::from("5.000").raw())),
);
assert_eq!(order.client_order_id(), client_order_id);
}
#[rstest]
fn test_queue_price_index_filters_other_prices() {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let (mut engine, _cache) = get_l3_queue_engine(instrument);
let target_price = Price::from("100.00");
let other_price = Price::from("101.00");
let target_id = ClientOrderId::from("O-QUEUE-TARGET");
let other_id = ClientOrderId::from("O-QUEUE-OTHER");
rest_l3_queue_order(&mut engine, target_price, 1, target_id);
rest_l3_queue_order(&mut engine, other_price, 2, other_id);
let indexed_ids = engine.take_queue_ids_at_price(target_price.raw());
assert_eq!(indexed_ids, vec![target_id]);
assert!(
engine
.queue_ids_by_price
.get(&other_price.raw())
.is_some_and(|ids| ids.contains(&other_id)),
);
}
fn rest_l3_queue_order(
engine: &mut OrderMatchingEngine,
price: Price,
sequence: u64,
client_order_id: ClientOrderId,
) -> OrderAny {
let instrument_id = engine.instrument.id();
let delta = OrderBookDelta::new(
instrument_id,
BookAction::Add,
BookOrder::new(OrderSide::Sell, price, Quantity::from("10.000"), sequence),
0,
sequence,
UnixNanos::from(sequence),
UnixNanos::from(sequence),
);
engine.process_order_book_delta(&delta).unwrap();
let mut order = OrderTestBuilder::new(OrderType::Limit)
.instrument_id(instrument_id)
.side(OrderSide::Sell)
.price(price)
.quantity(Quantity::from("5.000"))
.client_order_id(client_order_id)
.submit(true)
.build();
engine.process_order(&mut order, AccountId::from("SIM-001"));
order
}
#[derive(Debug, Clone, Copy)]
enum QueueEvent {
Add { id: OrderId, size: u64 },
Update { id: OrderId, size: u64 },
MoveAway { id: OrderId },
Delete { id: OrderId },
Trade { size: u64, aggressor: u8 },
AggregateCap { size: u64 },
AggregateDelete,
RestOrder,
}
fn granular_queue_event() -> impl Strategy<Value = QueueEvent> {
prop_oneof![
3 => (1u64..=6, 1u64..=9).prop_map(|(id, size)| QueueEvent::Add { id, size }),
3 => (1u64..=6, 1u64..=9).prop_map(|(id, size)| QueueEvent::Update { id, size }),
1 => (1u64..=6).prop_map(|id| QueueEvent::MoveAway { id }),
2 => (1u64..=6).prop_map(|id| QueueEvent::Delete { id }),
2 => Just(QueueEvent::RestOrder),
]
}
fn any_queue_event() -> impl Strategy<Value = QueueEvent> {
prop_oneof![
5 => granular_queue_event(),
3 => (1u64..=9, 0u8..3).prop_map(|(size, aggressor)| QueueEvent::Trade {
size,
aggressor,
}),
1 => (1u64..=9).prop_map(|size| QueueEvent::AggregateCap { size }),
1 => Just(QueueEvent::AggregateDelete),
]
}
struct L3QueueSim {
engine: OrderMatchingEngine,
account_id: AccountId,
live_main: HashMap<OrderId, u64>,
live_away: HashSet<OrderId>,
rest_snapshots: HashMap<ClientOrderId, HashSet<OrderId>>,
rested: usize,
sequence: u64,
}
impl L3QueueSim {
const MAIN_PRICE: &'static str = "100.00";
const AWAY_PRICE: &'static str = "101.00";
fn new() -> Self {
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let (engine, _cache) = get_l3_queue_engine(instrument);
Self {
engine,
account_id: AccountId::from("SIM-001"),
live_main: HashMap::new(),
live_away: HashSet::new(),
rest_snapshots: HashMap::new(),
rested: 0,
sequence: 0,
}
}
fn quantity(size: u64) -> Quantity {
Quantity::from(format!("{size}.000").as_str())
}
fn process_delta(
&mut self,
action: BookAction,
price: &str,
size: u64,
order_id: OrderId,
flags: u8,
) {
self.sequence += 1;
let delta = OrderBookDelta::new(
self.engine.instrument.id(),
action,
BookOrder::new(
OrderSide::Sell,
Price::from(price),
Self::quantity(size),
order_id,
),
flags,
self.sequence,
UnixNanos::from(self.sequence),
UnixNanos::from(self.sequence),
);
self.engine.process_order_book_delta(&delta).unwrap();
}
fn apply(&mut self, event: QueueEvent) {
match event {
QueueEvent::Add { id, size } => {
if self.live_main.contains_key(&id) || self.live_away.contains(&id) {
return;
}
self.process_delta(BookAction::Add, Self::MAIN_PRICE, size, id, 0);
self.live_main.insert(id, size);
}
QueueEvent::Update { id, size } => {
if !self.live_main.contains_key(&id) {
return;
}
self.process_delta(BookAction::Update, Self::MAIN_PRICE, size, id, 0);
self.live_main.insert(id, size);
}
QueueEvent::MoveAway { id } => {
let Some(size) = self.live_main.remove(&id) else {
return;
};
self.process_delta(BookAction::Update, Self::AWAY_PRICE, size, id, 0);
self.live_away.insert(id);
}
QueueEvent::Delete { id } => {
if let Some(size) = self.live_main.remove(&id) {
self.process_delta(BookAction::Delete, Self::MAIN_PRICE, size, id, 0);
} else if self.live_away.remove(&id) {
self.process_delta(BookAction::Delete, Self::AWAY_PRICE, 1, id, 0);
} else {
self.process_delta(BookAction::Delete, Self::MAIN_PRICE, 1, id, 0);
}
for snapshot_ids in self.rest_snapshots.values_mut() {
snapshot_ids.remove(&id);
}
}
QueueEvent::Trade { size, aggressor } => {
self.sequence += 1;
let aggressor_side = match aggressor {
0 => AggressorSide::Buy,
1 => AggressorSide::Sell,
_ => AggressorSide::NoAggressor,
};
let trade = TradeTick::new(
self.engine.instrument.id(),
Price::from(Self::MAIN_PRICE),
Self::quantity(size),
aggressor_side,
TradeId::new(format!("T-{}", self.sequence).as_str()),
UnixNanos::from(self.sequence),
UnixNanos::from(self.sequence),
);
self.engine.process_trade_tick(&trade);
}
QueueEvent::AggregateCap { size } => {
self.process_delta(
BookAction::Update,
Self::MAIN_PRICE,
size,
0,
RecordFlag::F_MBP as u8,
);
}
QueueEvent::AggregateDelete => {
self.process_delta(
BookAction::Delete,
Self::MAIN_PRICE,
1,
0,
RecordFlag::F_MBP as u8,
);
}
QueueEvent::RestOrder => {
if self.rested >= 3 {
return;
}
self.rested += 1;
let mut order = OrderTestBuilder::new(OrderType::Limit)
.instrument_id(self.engine.instrument.id())
.side(OrderSide::Sell)
.price(Price::from(Self::MAIN_PRICE))
.quantity(Self::quantity(5))
.client_order_id(ClientOrderId::from(
format!("O-PROP-{}", self.rested).as_str(),
))
.submit(true)
.build();
self.engine.process_order(&mut order, self.account_id);
assert!(
self.engine
.queue_ahead_orders
.contains_key(&order.client_order_id()),
"L3 snapshot must track the resting order",
);
self.rest_snapshots.insert(
order.client_order_id(),
self.live_main.keys().copied().collect(),
);
}
}
}
fn assert_tracked_orders_match_book(&self) {
let level: HashMap<OrderId, QuantityRaw> = self
.engine
.book
.get_orders_at_level(Price::from(Self::MAIN_PRICE), OrderSide::Buy)
.iter()
.map(|order| (order.order_id, order.size.raw()))
.collect();
for (client_order_id, orders_ahead) in &self.engine.queue_ahead_orders {
for (order_id, size_raw) in orders_ahead {
let book_size = level.get(order_id).copied().unwrap_or_else(|| {
panic!("tracked order {order_id} for {client_order_id} not in book level")
});
assert_eq!(
book_size, *size_raw,
"tracked size diverged from book for order {order_id}",
);
}
let tracked: HashSet<OrderId> = orders_ahead.keys().copied().collect();
let expected: HashSet<OrderId> = self.rest_snapshots[client_order_id]
.iter()
.filter(|id| self.live_main.contains_key(id))
.copied()
.collect();
assert_eq!(
tracked, expected,
"tracked set incomplete or stale for {client_order_id}",
);
}
}
}
#[rstest]
fn prop_test_l3_queue_tracking_stays_synced_with_counter() {
proptest!(|(events in prop::collection::vec(any_queue_event(), 1..=80))| {
let mut sim = L3QueueSim::new();
for event in events {
sim.apply(event);
assert_l3_queue_synced(&sim.engine);
}
});
}
#[rstest]
fn prop_test_l3_queue_tracking_mirrors_book_without_trades() {
proptest!(|(events in prop::collection::vec(granular_queue_event(), 1..=80))| {
let mut sim = L3QueueSim::new();
for event in events {
sim.apply(event);
assert_l3_queue_synced(&sim.engine);
sim.assert_tracked_orders_match_book();
}
});
}
#[rstest]
fn test_l3_queue_position_replay_databento_mbo_stays_synced() {
let json = include_str!("../../../../test_data/databento/esh4-glbx-mdp3-20231225.mbo.json");
let records: Vec<serde_json::Value> = serde_json::from_str(json).unwrap();
assert!(records.len() > 1000);
let instrument = InstrumentAny::FuturesContract(futures_contract_es(None, None));
let instrument_id = instrument.id();
let (mut engine, cache) = get_l3_queue_engine(instrument);
let account_id = AccountId::from("SIM-001");
let mut rested = 0usize;
let mut trades = 0usize;
for (index, record) in records.iter().enumerate() {
match record.get("type").and_then(serde_json::Value::as_str) {
Some("OrderBookDelta") => {
let mut delta: OrderBookDelta = serde_json::from_value(record.clone()).unwrap();
delta.instrument_id = instrument_id;
engine.process_order_book_delta(&delta).unwrap();
}
Some("TradeTick") => {
let mut trade: TradeTick = serde_json::from_value(record.clone()).unwrap();
trade.instrument_id = instrument_id;
engine.process_trade_tick(&trade);
trades += 1;
}
other => panic!("unexpected record type {other:?}"),
}
if index % 150 == 100 {
let (side, price) = if rested.is_multiple_of(2) {
(OrderSide::Sell, engine.book.best_ask_price())
} else {
(OrderSide::Buy, engine.book.best_bid_price())
};
if let Some(price) = price {
rested += 1;
let mut order = OrderTestBuilder::new(OrderType::Limit)
.instrument_id(instrument_id)
.side(side)
.price(price)
.quantity(Quantity::from("1"))
.client_order_id(ClientOrderId::from(format!("O-MBO-{rested}").as_str()))
.submit(true)
.build();
engine.process_order(&mut order, account_id);
let is_open = cache
.borrow()
.order(&order.client_order_id())
.is_some_and(|order| order.is_open());
if is_open {
assert!(
engine
.queue_ahead_orders
.contains_key(&order.client_order_id()),
"L3 snapshot must track the resting order",
);
}
}
}
assert_l3_queue_synced(&engine);
}
assert!(rested >= 5, "replay must exercise resting orders");
assert!(trades >= 50, "replay must exercise trade interleavings");
}
}