1use std::{cell::RefCell, collections::VecDeque, fmt::Debug, rc::Rc};
17
18use ahash::{AHashMap, AHashSet};
19use nautilus_common::{
20 cache::Cache,
21 clock::Clock,
22 logging::{CMD, EVT, RECV, SEND},
23 messages::{
24 data::{
25 DataCommand, SubscribeCommand, SubscribeQuotes, SubscribeTrades, UnsubscribeCommand,
26 UnsubscribeQuotes, UnsubscribeTrades,
27 },
28 execution::{
29 BatchModifyOrders, CancelAllOrders, CancelOrder, ModifyOrder, SubmitOrder,
30 SubmitOrderList, TradingCommand,
31 },
32 },
33 msgbus::{
34 self, MessagingSwitchboard, TypedHandler, TypedIntoHandler,
35 switchboard::{
36 get_event_order_topic, get_order_canceled_topic, get_quotes_topic, get_trades_topic,
37 },
38 },
39};
40use nautilus_core::{UUID4, WeakCell};
41use nautilus_model::{
42 data::{OrderBookDeltas, QuoteTick, TradeTick},
43 enums::{ContingencyType, OrderSide, OrderStatus, OrderType, TriggerType},
44 events::{OrderCanceled, OrderEmulated, OrderEventAny, OrderReleased, OrderUpdated},
45 identifiers::{ClientOrderId, ExecAlgorithmId, InstrumentId, PositionId, StrategyId},
46 instruments::Instrument,
47 orders::{LimitOrder, MarketOrder, Order, OrderAny},
48 types::{Price, Quantity},
49};
50use ustr::Ustr;
51
52use super::{PendingMessage, handlers::OrderEmulatorOnEventHandler};
53use crate::{
54 matching_core::{MatchAction, OrderMatchingCore, RestingOrder},
55 order_manager::{OrderManagerAction, manager::OrderManager},
56 trailing::trailing_stop_calculate,
57};
58
59pub struct OrderEmulator {
60 clock: Rc<RefCell<dyn Clock>>,
61 cache: Rc<RefCell<Cache>>,
62 manager: OrderManager,
63 matching_cores: AHashMap<InstrumentId, OrderMatchingCore>,
64 subscribed_quotes: AHashSet<InstrumentId>,
65 subscribed_trades: AHashSet<InstrumentId>,
66 subscribed_strategies: AHashSet<StrategyId>,
67 monitored_positions: AHashSet<PositionId>,
68 quote_tick_handler: Option<TypedHandler<QuoteTick>>,
69 trade_tick_handler: Option<TypedHandler<TradeTick>>,
70 quote_handlers: AHashMap<InstrumentId, TypedHandler<QuoteTick>>,
71 trade_handlers: AHashMap<InstrumentId, TypedHandler<TradeTick>>,
72 on_event_handler: Option<TypedHandler<OrderEventAny>>,
73 pending_messages: Rc<RefCell<VecDeque<PendingMessage>>>,
74}
75
76impl Debug for OrderEmulator {
77 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
78 f.debug_struct(stringify!(OrderEmulator))
79 .field("cores", &self.matching_cores.len())
80 .field("subscribed_quotes", &self.subscribed_quotes.len())
81 .field("subscribed_trades", &self.subscribed_trades.len())
82 .field("subscribed_strategies", &self.subscribed_strategies.len())
83 .finish()
84 }
85}
86
87impl OrderEmulator {
88 pub fn new(clock: Rc<RefCell<dyn Clock>>, cache: Rc<RefCell<Cache>>) -> Self {
89 let active_local = true;
90 let manager = OrderManager::new(clock.clone(), cache.clone(), active_local);
91
92 Self {
93 clock,
94 cache,
95 manager,
96 matching_cores: AHashMap::new(),
97 subscribed_quotes: AHashSet::new(),
98 subscribed_trades: AHashSet::new(),
99 subscribed_strategies: AHashSet::new(),
100 monitored_positions: AHashSet::new(),
101 quote_tick_handler: None,
102 trade_tick_handler: None,
103 quote_handlers: AHashMap::new(),
104 trade_handlers: AHashMap::new(),
105 on_event_handler: None,
106 pending_messages: Rc::new(RefCell::new(VecDeque::new())),
107 }
108 }
109
110 pub fn register_msgbus_handlers(emulator: &Rc<RefCell<Self>>) {
111 let weak = WeakCell::from(Rc::downgrade(emulator));
112 let pending_messages = emulator.borrow().pending_messages.clone();
113
114 let execute_weak = weak.clone();
115 let execute_pending = pending_messages.clone();
116 let execute_handler = TypedIntoHandler::from(move |cmd: TradingCommand| {
117 if let Some(emulator) = execute_weak.upgrade() {
118 match emulator.try_borrow_mut() {
119 Ok(mut emulator) => emulator.execute(cmd),
120 Err(_) => execute_pending
121 .borrow_mut()
122 .push_back(PendingMessage::Command(Box::new(cmd))),
123 }
124 }
125 });
126 msgbus::register_trading_command_endpoint(
127 MessagingSwitchboard::order_emulator_execute(),
128 execute_handler,
129 );
130
131 let quote_weak = weak.clone();
132 let quote_handler = TypedHandler::from(move |quote: &QuoteTick| {
133 if let Some(emulator) = quote_weak.upgrade() {
134 emulator.borrow_mut().on_quote_tick(*quote);
135 }
136 });
137
138 let trade_weak = weak.clone();
139 let trade_handler = TypedHandler::from(move |trade: &TradeTick| {
140 if let Some(emulator) = trade_weak.upgrade() {
141 emulator.borrow_mut().on_trade_tick(*trade);
142 }
143 });
144
145 let on_event_handler = TypedHandler::new(OrderEmulatorOnEventHandler::new(
146 Ustr::from(UUID4::new().as_str()),
147 weak,
148 WeakCell::from(Rc::downgrade(&pending_messages)),
149 ));
150
151 let mut emulator = emulator.borrow_mut();
152 emulator.quote_tick_handler = Some(quote_handler);
153 emulator.trade_tick_handler = Some(trade_handler);
154 emulator.on_event_handler = Some(on_event_handler);
155 }
156
157 pub fn set_on_event_handler(&mut self, handler: TypedHandler<OrderEventAny>) {
158 self.on_event_handler = Some(handler);
159 }
160
161 pub fn cache_submit_order_command(&mut self, command: SubmitOrder) {
163 self.manager.cache_submit_order_command(command);
164 }
165
166 fn subscribe_quotes_for_instrument(&mut self, instrument_id: InstrumentId) -> bool {
168 if self.quote_handlers.contains_key(&instrument_id) {
169 log::warn!("OrderEmulator attempted duplicate quote subscription for {instrument_id}");
170 return true;
171 }
172
173 let Some(handler) = self.quote_tick_handler.clone() else {
174 log::warn!("Cannot subscribe to quotes: msgbus handlers not registered");
175 return false;
176 };
177
178 let topic = get_quotes_topic(instrument_id);
179 self.quote_handlers.insert(instrument_id, handler.clone());
180 msgbus::subscribe_quotes(topic.into(), handler, None);
181 self.send_data_command(DataCommand::Subscribe(SubscribeCommand::Quotes(
182 SubscribeQuotes::new(
183 instrument_id,
184 None,
185 Some(instrument_id.venue),
186 UUID4::new(),
187 self.clock.borrow().timestamp_ns(),
188 None,
189 None,
190 ),
191 )));
192
193 true
194 }
195
196 fn subscribe_trades_for_instrument(&mut self, instrument_id: InstrumentId) -> bool {
198 if self.trade_handlers.contains_key(&instrument_id) {
199 log::warn!("OrderEmulator attempted duplicate trade subscription for {instrument_id}");
200 return true;
201 }
202
203 let Some(handler) = self.trade_tick_handler.clone() else {
204 log::warn!("Cannot subscribe to trades: msgbus handlers not registered");
205 return false;
206 };
207
208 let topic = get_trades_topic(instrument_id);
209 self.trade_handlers.insert(instrument_id, handler.clone());
210 msgbus::subscribe_trades(topic.into(), handler, None);
211 self.send_data_command(DataCommand::Subscribe(SubscribeCommand::Trades(
212 SubscribeTrades::new(
213 instrument_id,
214 None,
215 Some(instrument_id.venue),
216 UUID4::new(),
217 self.clock.borrow().timestamp_ns(),
218 None,
219 None,
220 ),
221 )));
222
223 true
224 }
225
226 #[must_use]
227 pub fn subscribed_quotes(&self) -> Vec<InstrumentId> {
228 let mut quotes: Vec<InstrumentId> = self.subscribed_quotes.iter().copied().collect();
229 quotes.sort();
230 quotes
231 }
232
233 #[must_use]
234 pub fn subscribed_trades(&self) -> Vec<InstrumentId> {
235 let mut trades: Vec<_> = self.subscribed_trades.iter().copied().collect();
236 trades.sort();
237 trades
238 }
239
240 #[must_use]
241 pub fn subscribed_strategy_count(&self) -> usize {
242 self.subscribed_strategies.len()
243 }
244
245 #[must_use]
246 pub fn monitored_position_count(&self) -> usize {
247 self.monitored_positions.len()
248 }
249
250 #[must_use]
251 pub fn get_submit_order_commands(&self) -> AHashMap<ClientOrderId, SubmitOrder> {
252 self.manager.get_submit_order_commands()
253 }
254
255 #[must_use]
256 pub fn get_matching_core(&self, instrument_id: &InstrumentId) -> Option<OrderMatchingCore> {
257 self.matching_cores.get(instrument_id).cloned()
258 }
259
260 pub fn start(&mut self) {
261 if let Err(e) = self.on_start() {
262 log::error!("{e}");
263 }
264
265 log::info!("Started");
266 }
267
268 pub fn stop(&self) {
269 self.on_stop();
270
271 log::info!("Stopped");
272 }
273
274 pub fn reset(&mut self) {
275 self.on_reset();
276
277 log::info!("Reset");
278 }
279
280 pub fn dispose(&mut self) {
281 self.on_dispose();
282
283 log::info!("Disposed");
284 }
285
286 pub fn on_start(&mut self) -> anyhow::Result<()> {
292 let emulated_orders: Vec<OrderAny> = self
293 .cache
294 .borrow()
295 .orders_emulated(None, None, None, None, None)
296 .into_iter()
297 .map(|o| o.clone())
298 .collect();
299
300 if emulated_orders.is_empty() {
301 log::debug!("No emulated orders to reactivate");
302 return Ok(());
303 }
304
305 for order in emulated_orders {
306 if !matches!(
307 order.status(),
308 OrderStatus::Initialized | OrderStatus::Emulated
309 ) {
310 continue; }
312
313 if let Some(parent_order_id) = &order.parent_order_id() {
314 let parent_order = if let Some(order) = self.cache.borrow().order(parent_order_id) {
315 order.clone()
316 } else {
317 log::error!("Cannot handle order: parent {parent_order_id} not found");
318 continue;
319 };
320
321 let is_position_closed = parent_order
322 .position_id()
323 .is_none_or(|id| self.cache.borrow().is_position_closed(&id));
324 if parent_order.is_closed() && is_position_closed {
325 let actions = self.manager.cancel_order(&order);
326 self.dispatch_manager_actions(actions);
327 continue; }
329
330 if parent_order.contingency_type() == Some(ContingencyType::Oto)
331 && (parent_order.is_active_local()
332 || parent_order.filled_qty() == Quantity::zero(0))
333 {
334 continue; }
336 }
337
338 let position_id = self
339 .cache
340 .borrow()
341 .position_id(&order.client_order_id())
342 .copied();
343 let client_id = self
344 .cache
345 .borrow()
346 .client_id(&order.client_order_id())
347 .copied();
348
349 let command = SubmitOrder::new(
350 order.trader_id(),
351 client_id,
352 order.strategy_id(),
353 order.instrument_id(),
354 order.client_order_id(),
355 order.init_event().clone(),
356 order.exec_algorithm_id(),
357 position_id,
358 None, UUID4::new(),
360 self.clock.borrow().timestamp_ns(),
361 None, );
363
364 self.manager.cache_submit_order_command(command.clone());
365 self.handle_submit_order(&command);
366 }
367
368 self.drain_pending_messages();
369
370 Ok(())
371 }
372
373 pub fn on_event(&mut self, event: &OrderEventAny) {
374 log::info!("{RECV}{EVT} {event}");
375
376 let actions = self.manager.handle_event(event);
377 self.dispatch_manager_actions(actions);
378
379 if let Some(order) = self.cache.borrow().order(&event.client_order_id())
380 && order.is_closed()
381 && let Some(matching_core) = self.matching_cores.get_mut(&order.instrument_id())
382 && let Err(e) = matching_core.delete_order(event.client_order_id())
383 {
384 log::debug!("Error deleting order: {e}");
385 }
386 self.drain_pending_messages();
389 }
390
391 fn drain_pending_messages(&mut self) {
392 loop {
393 let message = self.pending_messages.borrow_mut().pop_front();
394 match message {
395 Some(PendingMessage::Command(command)) => self.execute(*command),
396 Some(PendingMessage::Event(event)) => self.on_event(&event),
397 None => break,
398 }
399 }
400 }
401
402 pub const fn on_stop(&self) {}
403
404 pub fn on_reset(&mut self) {
405 self.manager.reset();
406 self.pending_messages.borrow_mut().clear();
407 self.matching_cores.clear();
408 self.unsubscribe_all_market_data();
409 self.unsubscribe_strategy_order_events();
410 self.monitored_positions.clear();
411 }
412
413 pub fn on_dispose(&mut self) {
414 self.on_reset();
415 }
416
417 fn unsubscribe_all_market_data(&mut self) {
418 let mut quote_instrument_ids: Vec<_> = self.subscribed_quotes.drain().collect();
421 quote_instrument_ids.sort();
422
423 for instrument_id in quote_instrument_ids {
424 self.unsubscribe_quotes_for_instrument(instrument_id);
425 }
426
427 let mut trade_instrument_ids: Vec<_> = self.subscribed_trades.drain().collect();
428 trade_instrument_ids.sort();
429
430 for instrument_id in trade_instrument_ids {
431 self.unsubscribe_trades_for_instrument(instrument_id);
432 }
433 }
434
435 fn unsubscribe_quotes_for_instrument(&mut self, instrument_id: InstrumentId) {
436 if let Some(handler) = self.quote_handlers.remove(&instrument_id) {
437 let topic = get_quotes_topic(instrument_id);
438 msgbus::unsubscribe_quotes(topic.into(), &handler);
439 }
440
441 self.send_data_command(DataCommand::Unsubscribe(UnsubscribeCommand::Quotes(
442 UnsubscribeQuotes::new(
443 instrument_id,
444 None,
445 Some(instrument_id.venue),
446 UUID4::new(),
447 self.clock.borrow().timestamp_ns(),
448 None,
449 None,
450 ),
451 )));
452 }
453
454 fn unsubscribe_trades_for_instrument(&mut self, instrument_id: InstrumentId) {
455 if let Some(handler) = self.trade_handlers.remove(&instrument_id) {
456 let topic = get_trades_topic(instrument_id);
457 msgbus::unsubscribe_trades(topic.into(), &handler);
458 }
459
460 self.send_data_command(DataCommand::Unsubscribe(UnsubscribeCommand::Trades(
461 UnsubscribeTrades::new(
462 instrument_id,
463 None,
464 Some(instrument_id.venue),
465 UUID4::new(),
466 self.clock.borrow().timestamp_ns(),
467 None,
468 None,
469 ),
470 )));
471 }
472
473 fn unsubscribe_strategy_order_events(&mut self) {
474 let mut strategy_ids: Vec<_> = self.subscribed_strategies.drain().collect();
475 strategy_ids.sort();
476 let Some(handler) = &self.on_event_handler else {
477 return;
478 };
479
480 for strategy_id in strategy_ids {
481 msgbus::unsubscribe_order_events(format!("events.order.{strategy_id}").into(), handler);
482 }
483 }
484
485 pub fn execute(&mut self, command: TradingCommand) {
486 log::info!("{RECV}{CMD} {command}");
487
488 match command {
489 TradingCommand::SubmitOrder(command) => self.handle_submit_order(&command),
490 TradingCommand::SubmitOrderList(ref command) => self.handle_submit_order_list(command),
491 TradingCommand::ModifyOrder(ref command) => self.handle_modify_order(command),
492 TradingCommand::ModifyOrders(ref command) => self.handle_batch_modify_orders(command),
493 TradingCommand::CancelOrder(command) => self.handle_cancel_order(command),
494 TradingCommand::CancelAllOrders(ref command) => self.handle_cancel_all_orders(command),
495 _ => log::error!("Cannot handle command: unrecognized {command:?}"),
496 }
497
498 self.drain_pending_messages();
499 }
500
501 fn dispatch_manager_actions(&mut self, actions: Vec<OrderManagerAction>) {
502 for action in actions {
503 self.dispatch_manager_action(action);
504 }
505 }
506
507 fn dispatch_manager_action(&mut self, action: OrderManagerAction) {
508 match action {
509 OrderManagerAction::PublishInitialized(event) => publish_order_event(&event),
510 OrderManagerAction::SubmitToEmulator(command) => self.handle_submit_order(&command),
511 OrderManagerAction::SubmitToRisk(command) => {
512 self.send_risk_command(TradingCommand::SubmitOrder(command));
513 }
514 OrderManagerAction::SubmitToAlgorithm {
515 command,
516 exec_algorithm_id,
517 } => self.send_algo_command(command, exec_algorithm_id),
518 OrderManagerAction::CancelLocal(order) => self.cancel_order(&order),
519 OrderManagerAction::ModifyLocalQuantity {
520 mut order,
521 quantity,
522 } => {
523 self.update_order(&mut order, quantity);
524 }
525 }
526 }
527
528 fn create_matching_core(
529 &mut self,
530 instrument_id: InstrumentId,
531 price_increment: Price,
532 ) -> OrderMatchingCore {
533 let matching_core = OrderMatchingCore::new(instrument_id, price_increment);
534 self.matching_cores
535 .insert(instrument_id, matching_core.clone());
536 log::info!("Creating matching core for {instrument_id:?}");
537 matching_core
538 }
539
540 pub fn handle_submit_order(&mut self, command: &SubmitOrder) {
544 let client_order_id = command.client_order_id;
545
546 let mut order = self
547 .cache
548 .borrow()
549 .order(&client_order_id)
550 .map(|o| o.clone())
551 .expect("order must exist in cache");
552
553 let emulation_trigger = order.emulation_trigger();
554
555 if !self
556 .manager
557 .get_submit_order_commands()
558 .contains_key(&client_order_id)
559 {
560 self.manager.cache_submit_order_command(command.clone());
561 }
562
563 assert!(
564 emulation_trigger.is_some(),
565 "order.emulation_trigger must be set"
566 );
567
568 if !matches!(
569 emulation_trigger,
570 Some(TriggerType::Default | TriggerType::BidAsk | TriggerType::LastPrice)
571 ) {
572 log::error!("Cannot emulate order: `TriggerType` {emulation_trigger:?} not supported");
573 let actions = self.manager.cancel_order(&order);
574 self.dispatch_manager_actions(actions);
575 return;
576 }
577 let strategy_id = command.strategy_id;
578 let position_id = command.position_id;
579
580 let trigger_instrument_id = order
582 .trigger_instrument_id()
583 .unwrap_or_else(|| order.instrument_id());
584
585 let matching_core = self.matching_cores.get(&trigger_instrument_id).cloned();
586
587 let mut matching_core = if let Some(core) = matching_core {
588 core
589 } else {
590 let (instrument_id, price_increment) = if trigger_instrument_id.is_synthetic() {
592 let synthetic = self
593 .cache
594 .borrow()
595 .try_synthetic(&trigger_instrument_id)
596 .cloned();
597
598 match synthetic {
599 Ok(synthetic) => (synthetic.id, synthetic.price_increment),
600 Err(e) => {
601 log::error!("Cannot emulate order: {e}");
602 let actions = self.manager.cancel_order(&order);
603 self.dispatch_manager_actions(actions);
604 return;
605 }
606 }
607 } else {
608 let instrument = self
609 .cache
610 .borrow()
611 .instrument(&trigger_instrument_id)
612 .cloned();
613
614 if let Some(instrument) = instrument {
615 (instrument.id(), instrument.price_increment())
616 } else {
617 log::error!(
618 "Cannot emulate order: no instrument {trigger_instrument_id} for trigger"
619 );
620 let actions = self.manager.cancel_order(&order);
621 self.dispatch_manager_actions(actions);
622 return;
623 }
624 };
625
626 self.create_matching_core(instrument_id, price_increment)
627 };
628
629 if matches!(
631 order.order_type(),
632 OrderType::TrailingStopMarket | OrderType::TrailingStopLimit
633 ) {
634 self.update_trailing_stop_order(&mut order);
635 if order.trigger_price().is_none() && is_order_activated(&order) {
636 log::error!(
637 "Cannot handle trailing stop order with no trigger_price and no market updates"
638 );
639
640 let actions = self.manager.cancel_order(&order);
641 self.dispatch_manager_actions(actions);
642 return;
643 }
644 }
646
647 self.check_monitoring(strategy_id, position_id);
648
649 let is_activated = is_order_activated(&order);
651 let match_info = RestingOrder::new(
652 order.client_order_id(),
653 order.order_side(),
654 order.order_type(),
655 if is_activated {
656 order.trigger_price()
657 } else {
658 None
659 },
660 if is_activated { order.price() } else { None },
661 is_activated,
662 );
663
664 if let Some(action) = matching_core.match_order(&match_info) {
665 self.dispatch_match_action(action);
666 }
667
668 match emulation_trigger.unwrap() {
670 TriggerType::Default | TriggerType::BidAsk => {
671 if !self.subscribed_quotes.contains(&trigger_instrument_id)
672 && self.subscribe_quotes_for_instrument(trigger_instrument_id)
673 {
674 self.subscribed_quotes.insert(trigger_instrument_id);
675 }
676 }
677 TriggerType::LastPrice => {
678 if !self.subscribed_trades.contains(&trigger_instrument_id)
679 && self.subscribe_trades_for_instrument(trigger_instrument_id)
680 {
681 self.subscribed_trades.insert(trigger_instrument_id);
682 }
683 }
684 _ => {
685 log::error!("Invalid TriggerType: {emulation_trigger:?}");
686 return;
687 }
688 }
689
690 if !self
692 .manager
693 .get_submit_order_commands()
694 .contains_key(&order.client_order_id())
695 {
696 return; }
698
699 matching_core.add_order(match_info);
701
702 if order.status() == OrderStatus::Initialized {
704 let event = OrderEmulated::new(
705 order.trader_id(),
706 order.strategy_id(),
707 order.instrument_id(),
708 order.client_order_id(),
709 UUID4::new(),
710 self.clock.borrow().timestamp_ns(),
711 self.clock.borrow().timestamp_ns(),
712 );
713
714 let event = OrderEventAny::Emulated(event);
715
716 order = match self.cache.borrow_mut().update_order(&event) {
717 Ok(order) => order,
718 Err(e) => {
719 log::error!("Cannot apply order event: {e:?}");
720 return;
721 }
722 };
723
724 self.send_risk_event(event.clone());
725
726 msgbus::publish_order_event(
727 format!("events.order.{}", order.strategy_id()).into(),
728 &event,
729 );
730 }
731
732 self.matching_cores
734 .insert(trigger_instrument_id, matching_core);
735
736 log::info!("Emulating {order}");
737 }
738
739 fn handle_submit_order_list(&mut self, command: &SubmitOrderList) {
740 self.check_monitoring(command.strategy_id, command.position_id);
741
742 let orders: Vec<OrderAny> = self
743 .cache
744 .borrow()
745 .orders_for_ids(&command.order_list.client_order_ids, &command);
746
747 for order in &orders {
748 if let Some(parent_order_id) = order.parent_order_id() {
749 let cache = self.cache.borrow();
750 let parent_order = if let Some(parent_order) = cache.order(&parent_order_id) {
751 parent_order
752 } else {
753 log::error!("Parent order for {} not found", order.client_order_id());
754 continue;
755 };
756
757 if parent_order.contingency_type() == Some(ContingencyType::Oto) {
758 continue; }
760 }
761
762 match self.manager.create_new_submit_order(
763 order,
764 command.position_id,
765 command.client_id,
766 command.correlation_id,
767 ) {
768 Ok(actions) => self.dispatch_manager_actions(actions),
769 Err(e) => log::error!("Error creating new submit order: {e}"),
770 }
771 }
772 }
773
774 fn handle_modify_order(&mut self, command: &ModifyOrder) {
775 let order = self
776 .cache
777 .borrow()
778 .order(&command.client_order_id)
779 .map(|order| order.clone());
780
781 if let Some(order) = order {
782 let price = match command.price {
783 Some(price) => Some(price),
784 None => order.price(),
785 };
786
787 let trigger_price = match command.trigger_price {
788 Some(trigger_price) => Some(trigger_price),
789 None => order.trigger_price(),
790 };
791
792 let ts_now = self.clock.borrow().timestamp_ns();
793 let event = OrderUpdated::new(
794 order.trader_id(),
795 order.strategy_id(),
796 order.instrument_id(),
797 order.client_order_id(),
798 command.quantity.unwrap_or(order.quantity()),
799 UUID4::new(),
800 ts_now,
801 ts_now,
802 false,
803 order.venue_order_id(),
804 order.account_id(),
805 price,
806 trigger_price,
807 None,
808 order.is_quote_quantity(),
809 );
810
811 let event = OrderEventAny::Updated(event);
812 self.send_exec_event(event.clone());
813
814 let order = self
816 .cache
817 .borrow()
818 .order(&command.client_order_id)
819 .filter(|order| order.last_event() == &event)
820 .map(|order| order.clone());
821
822 let Some(order) = order else {
823 return;
824 };
825
826 if !self
827 .manager
828 .get_submit_order_commands()
829 .contains_key(&command.client_order_id)
830 {
831 return;
832 }
833
834 let trigger_instrument_id = order
835 .trigger_instrument_id()
836 .unwrap_or_else(|| order.instrument_id());
837
838 let action = if let Some(matching_core) =
839 self.matching_cores.get_mut(&trigger_instrument_id)
840 {
841 if let Err(e) = matching_core.delete_order(order.client_order_id()) {
842 log::debug!("Cannot update order match info: {e:?}");
843 }
844
845 let is_activated = is_order_activated(&order);
846 let match_info = RestingOrder::new(
847 order.client_order_id(),
848 order.order_side(),
849 order.order_type(),
850 if is_activated {
851 order.trigger_price()
852 } else {
853 None
854 },
855 if is_activated { order.price() } else { None },
856 is_activated,
857 );
858 let action = matching_core.match_order(&match_info);
859 if action.is_none() {
860 matching_core.add_order(match_info);
861 }
862 action
863 } else {
864 log::error!(
865 "Cannot handle `ModifyOrder`: no matching core for trigger instrument {trigger_instrument_id}"
866 );
867 return;
868 };
869
870 if let Some(action) = action {
871 self.dispatch_match_action(action);
872 }
873 } else {
874 log::error!("Cannot modify order: {} not found", command.client_order_id);
875 }
876 }
877
878 fn handle_batch_modify_orders(&mut self, command: &BatchModifyOrders) {
879 for modify in &command.modifies {
880 self.handle_modify_order(modify);
881 }
882 }
883
884 pub fn handle_cancel_order(&mut self, command: CancelOrder) {
885 let order = if let Some(order) = self.cache.borrow().order(&command.client_order_id) {
886 order.clone()
887 } else {
888 log::error!("Cannot cancel order: {} not found", command.client_order_id);
889 return;
890 };
891
892 let trigger_instrument_id = order
893 .trigger_instrument_id()
894 .unwrap_or_else(|| order.instrument_id());
895
896 let matching_core = if let Some(core) = self.matching_cores.get(&trigger_instrument_id) {
897 core
898 } else {
899 let actions = self.manager.cancel_order(&order);
900 self.dispatch_manager_actions(actions);
901 return;
902 };
903
904 if !matching_core.order_exists(order.client_order_id())
905 && order.is_open()
906 && !order.is_pending_cancel()
907 {
908 self.send_exec_command(TradingCommand::CancelOrder(command));
910 } else {
911 let actions = self.manager.cancel_order(&order);
912 self.dispatch_manager_actions(actions);
913 }
914 }
915
916 fn handle_cancel_all_orders(&mut self, command: &CancelAllOrders) {
917 let mut ids_to_cancel: Vec<ClientOrderId> = self
918 .matching_cores
919 .values()
920 .flat_map(OrderMatchingCore::iter_orders)
921 .map(|order| order.client_order_id)
922 .collect();
923 ids_to_cancel.sort_unstable();
924 ids_to_cancel.dedup();
925
926 {
927 let cache = self.cache.borrow();
928 ids_to_cancel.retain(|client_order_id| {
929 let Some(order) = cache.order(client_order_id) else {
930 return false;
931 };
932
933 order.instrument_id() == command.instrument_id
934 && command
935 .order_side
936 .is_none_or(|side| order.order_side() == side)
937 && command.client_id.is_none_or(|client_id| {
938 cache
939 .client_id(client_order_id)
940 .is_some_and(|order_client_id| *order_client_id == client_id)
941 })
942 });
943 }
944
945 for id in ids_to_cancel {
946 let Some(order) = self.cache.borrow().order_owned(&id) else {
947 continue;
948 };
949 let actions = self.manager.cancel_order(&order);
950 self.dispatch_manager_actions(actions);
951 }
952 }
953
954 pub fn update_order(&mut self, order: &mut OrderAny, new_quantity: Quantity) {
955 log::info!(
956 "Updating order {} quantity to {new_quantity}",
957 order.client_order_id(),
958 );
959
960 let ts_now = self.clock.borrow().timestamp_ns();
961 let event = OrderUpdated::new(
962 order.trader_id(),
963 order.strategy_id(),
964 order.instrument_id(),
965 order.client_order_id(),
966 new_quantity,
967 UUID4::new(),
968 ts_now,
969 ts_now,
970 false,
971 None,
972 order.account_id(),
973 None,
974 None,
975 None,
976 order.is_quote_quantity(),
977 );
978
979 let event = OrderEventAny::Updated(event);
980
981 *order = match self.cache.borrow_mut().update_order(&event) {
982 Ok(order) => order,
983 Err(e) => {
984 log::error!("Cannot apply order event: {e:?}");
985 return;
986 }
987 };
988
989 self.send_risk_event(event);
990 }
991
992 pub fn on_order_book_deltas(&mut self, deltas: &OrderBookDeltas) {
993 log::debug!("Processing {deltas:?}");
994
995 let instrument_id = &deltas.instrument_id;
996 if let Some(matching_core) = self.matching_cores.get_mut(instrument_id) {
997 if let Some(book) = self.cache.borrow().order_book(instrument_id) {
998 let best_bid = book.best_bid_price();
999 let best_ask = book.best_ask_price();
1000
1001 if let Some(best_bid) = best_bid {
1002 matching_core.set_bid_raw(best_bid);
1003 }
1004
1005 if let Some(best_ask) = best_ask {
1006 matching_core.set_ask_raw(best_ask);
1007 }
1008 } else {
1009 log::error!(
1010 "Cannot handle `OrderBookDeltas`: no book being maintained for {}",
1011 deltas.instrument_id
1012 );
1013 }
1014
1015 self.iterate_orders(instrument_id);
1016 } else {
1017 log::error!(
1018 "Cannot handle `OrderBookDeltas`: no matching core for instrument {}",
1019 deltas.instrument_id
1020 );
1021 }
1022 }
1023
1024 pub fn on_quote_tick(&mut self, quote: QuoteTick) {
1025 log::debug!("Processing {quote}:?");
1026
1027 let instrument_id = "e.instrument_id;
1028 if let Some(matching_core) = self.matching_cores.get_mut(instrument_id) {
1029 matching_core.set_bid_raw(quote.bid_price);
1030 matching_core.set_ask_raw(quote.ask_price);
1031
1032 self.iterate_orders(instrument_id);
1033 } else {
1034 log::error!(
1035 "Cannot handle `QuoteTick`: no matching core for instrument {}",
1036 quote.instrument_id
1037 );
1038 }
1039 }
1040
1041 pub fn on_trade_tick(&mut self, trade: TradeTick) {
1042 log::debug!("Processing {trade:?}");
1043
1044 let instrument_id = &trade.instrument_id;
1045 if let Some(matching_core) = self.matching_cores.get_mut(instrument_id) {
1046 matching_core.set_last_raw(trade.price);
1047
1048 if !self.subscribed_quotes.contains(instrument_id) {
1049 matching_core.set_bid_raw(trade.price);
1050 matching_core.set_ask_raw(trade.price);
1051 }
1052
1053 self.iterate_orders(instrument_id);
1054 } else {
1055 log::error!(
1056 "Cannot handle `TradeTick`: no matching core for instrument {}",
1057 trade.instrument_id
1058 );
1059 }
1060 }
1061
1062 fn iterate_orders(&mut self, instrument_id: &InstrumentId) {
1063 let bid_actions = if let Some(matching_core) = self.matching_cores.get_mut(instrument_id) {
1066 matching_core.iterate_bids()
1067 } else {
1068 log::error!("Cannot iterate orders: no matching core for instrument {instrument_id}");
1069 return;
1070 };
1071
1072 for action in bid_actions {
1073 self.dispatch_match_action(action);
1074 }
1075
1076 let ask_actions = if let Some(matching_core) = self.matching_cores.get_mut(instrument_id) {
1077 matching_core.iterate_asks()
1078 } else {
1079 return;
1080 };
1081
1082 for action in ask_actions {
1083 self.dispatch_match_action(action);
1084 }
1085
1086 let orders = if let Some(matching_core) = self.matching_cores.get(instrument_id) {
1088 matching_core.get_orders()
1089 } else {
1090 return;
1091 };
1092
1093 for match_info in orders {
1094 if !matches!(
1095 match_info.order_type,
1096 OrderType::TrailingStopMarket | OrderType::TrailingStopLimit
1097 ) {
1098 continue;
1099 }
1100
1101 let mut order = match self
1102 .cache
1103 .borrow()
1104 .order(&match_info.client_order_id)
1105 .map(|o| o.clone())
1106 {
1107 Some(order) => order,
1108 None => continue,
1109 };
1110
1111 if order.is_closed() {
1112 continue;
1113 }
1114
1115 self.update_trailing_stop_order(&mut order);
1116 }
1117
1118 self.drain_pending_messages();
1119 }
1120
1121 fn dispatch_match_action(&mut self, action: MatchAction) {
1122 match action {
1123 MatchAction::FillLimit(id) => self.fill_limit_order(id),
1124 MatchAction::TriggerStop(id) => self.trigger_stop_order(id),
1125 }
1126 }
1127
1128 pub fn cancel_order(&mut self, order: &OrderAny) {
1129 log::info!("Canceling order {}", order.client_order_id());
1130
1131 let mut order = order.clone();
1132 order.set_emulation_trigger(None);
1133
1134 let trigger_instrument_id = order
1135 .trigger_instrument_id()
1136 .unwrap_or(order.instrument_id());
1137
1138 if let Some(matching_core) = self.matching_cores.get_mut(&trigger_instrument_id)
1139 && let Err(e) = matching_core.delete_order(order.client_order_id())
1140 {
1141 log::debug!("Cannot delete order: {e:?}");
1142 }
1143
1144 self.manager
1145 .pop_submit_order_command(order.client_order_id());
1146
1147 self.cache
1148 .borrow_mut()
1149 .update_order_pending_cancel_local(&order);
1150
1151 let ts_now = self.clock.borrow().timestamp_ns();
1152 let event = OrderCanceled::new(
1153 order.trader_id(),
1154 order.strategy_id(),
1155 order.instrument_id(),
1156 order.client_order_id(),
1157 UUID4::new(),
1158 ts_now,
1159 ts_now,
1160 false,
1161 order.venue_order_id(),
1162 order.account_id(),
1163 );
1164
1165 let event = OrderEventAny::Canceled(event);
1166 if let Err(e) = self.cache.borrow_mut().update_order(&event) {
1167 log::error!("Failed to apply order event: {e}");
1168 return;
1169 }
1170
1171 self.send_portfolio_order_event(event.clone());
1172 publish_order_event(&event);
1173 }
1174
1175 fn check_monitoring(&mut self, strategy_id: StrategyId, position_id: Option<PositionId>) {
1176 if !self.subscribed_strategies.contains(&strategy_id) {
1177 if let Some(handler) = &self.on_event_handler {
1179 msgbus::subscribe_order_events(
1180 format!("events.order.{strategy_id}").into(),
1181 handler.clone(),
1182 None,
1183 );
1184 self.subscribed_strategies.insert(strategy_id);
1185 log::info!("Subscribed to strategy {strategy_id} order events");
1186 }
1187 }
1188
1189 if let Some(position_id) = position_id
1190 && !self.monitored_positions.contains(&position_id)
1191 {
1192 self.monitored_positions.insert(position_id);
1193 }
1194 }
1195
1196 fn validate_release(
1204 &self,
1205 order: &OrderAny,
1206 matching_core: &OrderMatchingCore,
1207 trigger_instrument_id: InstrumentId,
1208 ) -> Option<Price> {
1209 let released_price = match order.order_side() {
1210 OrderSide::Buy => matching_core.ask,
1211 OrderSide::Sell => matching_core.bid,
1212 };
1213
1214 if released_price.is_none() {
1215 log::warn!(
1216 "Cannot release order {} yet: no market data available for {trigger_instrument_id}, will retry on next update",
1217 order.client_order_id(),
1218 );
1219 return None;
1220 }
1221
1222 Some(released_price.unwrap())
1223 }
1224
1225 pub fn trigger_stop_order(&mut self, client_order_id: ClientOrderId) {
1229 let order = match self
1230 .cache
1231 .borrow()
1232 .order(&client_order_id)
1233 .map(|o| o.clone())
1234 {
1235 Some(order) => order,
1236 None => {
1237 log::error!(
1238 "Cannot trigger stop order: order {client_order_id} not found in cache"
1239 );
1240 return;
1241 }
1242 };
1243
1244 match order.order_type() {
1245 OrderType::StopLimit | OrderType::LimitIfTouched | OrderType::TrailingStopLimit => {
1246 self.fill_limit_order(client_order_id);
1247 }
1248 OrderType::Market
1249 | OrderType::MarketIfTouched
1250 | OrderType::StopMarket
1251 | OrderType::TrailingStopMarket => self.fill_market_order(client_order_id),
1252 _ => panic!("invalid `OrderType`, was {}", order.order_type()),
1253 }
1254 }
1255
1256 pub fn fill_limit_order(&mut self, client_order_id: ClientOrderId) {
1260 let order = match self
1261 .cache
1262 .borrow()
1263 .order(&client_order_id)
1264 .map(|o| o.clone())
1265 {
1266 Some(order) => order,
1267 None => {
1268 log::error!("Cannot fill limit order: order {client_order_id} not found in cache");
1269 return;
1270 }
1271 };
1272
1273 if matches!(order.order_type(), OrderType::Limit) {
1274 self.fill_market_order(client_order_id);
1275 return;
1276 }
1277
1278 let trigger_instrument_id = order
1279 .trigger_instrument_id()
1280 .unwrap_or(order.instrument_id());
1281
1282 let matching_core = match self.matching_cores.get(&trigger_instrument_id) {
1283 Some(core) => core,
1284 None => {
1285 log::error!(
1286 "Cannot fill limit order: no matching core for instrument {trigger_instrument_id}"
1287 );
1288 return; }
1290 };
1291
1292 let released_price =
1293 match self.validate_release(&order, matching_core, trigger_instrument_id) {
1294 Some(price) => price,
1295 None => return, };
1297
1298 let command = match self
1299 .manager
1300 .pop_submit_order_command(order.client_order_id())
1301 {
1302 Some(command) => command,
1303 None => return, };
1305
1306 if let Some(matching_core) = self.matching_cores.get_mut(&trigger_instrument_id) {
1307 if let Err(e) = matching_core.delete_order(client_order_id) {
1308 log::debug!("Error deleting order: {e:?}");
1309 }
1310
1311 let mut transformed = if let Ok(transformed) = LimitOrder::new_checked(
1313 order.trader_id(),
1314 order.strategy_id(),
1315 order.instrument_id(),
1316 order.client_order_id(),
1317 order.order_side(),
1318 order.quantity(),
1319 order.price().unwrap(),
1320 order.time_in_force(),
1321 order.expire_time(),
1322 order.is_post_only(),
1323 order.is_reduce_only(),
1324 order.is_quote_quantity(),
1325 order.display_qty(),
1326 None,
1327 Some(trigger_instrument_id),
1328 order.contingency_type(),
1329 order.order_list_id(),
1330 order.linked_order_ids().map(Vec::from),
1331 order.parent_order_id(),
1332 order.exec_algorithm_id(),
1333 order.exec_algorithm_params().cloned(),
1334 order.exec_spawn_id(),
1335 order.tags().map(Vec::from),
1336 UUID4::new(),
1337 self.clock.borrow().timestamp_ns(),
1338 ) {
1339 transformed
1340 } else {
1341 log::error!("Cannot create limit order");
1342 return;
1343 };
1344 transformed.liquidity_side = order.liquidity_side();
1345
1346 let original_events = order.events();
1353
1354 for event in original_events.into_iter().rev() {
1357 transformed.events.insert(0, event.clone());
1358 }
1359
1360 let add_result = {
1361 let mut cache = self.cache.borrow_mut();
1362 cache.add_order(
1363 OrderAny::Limit(transformed.clone()),
1364 command.position_id,
1365 command.client_id,
1366 true,
1367 )
1368 };
1369
1370 if let Err(e) = add_result {
1371 log::error!("Failed to add order: {e}");
1372 } else {
1373 msgbus::publish_order_event(
1374 format!("events.order.{}", order.strategy_id()).into(),
1375 transformed.last_event(),
1376 );
1377 }
1378
1379 let event = OrderReleased::new(
1380 order.trader_id(),
1381 order.strategy_id(),
1382 order.instrument_id(),
1383 order.client_order_id(),
1384 released_price,
1385 UUID4::new(),
1386 self.clock.borrow().timestamp_ns(),
1387 self.clock.borrow().timestamp_ns(),
1388 );
1389
1390 let event = OrderEventAny::Released(event);
1391
1392 let transformed = match self.cache.borrow_mut().update_order(&event) {
1393 Ok(order) => order,
1394 Err(e) => {
1395 log::error!("Failed to apply order event: {e}");
1396 return;
1397 }
1398 };
1399
1400 self.send_risk_event(event.clone());
1401
1402 log::info!("Releasing order {}", order.client_order_id());
1403
1404 msgbus::publish_order_event(
1406 format!("events.order.{}", transformed.strategy_id()).into(),
1407 &event,
1408 );
1409
1410 if let Some(exec_algorithm_id) = order.exec_algorithm_id() {
1411 self.send_algo_command(command, exec_algorithm_id);
1412 } else {
1413 self.send_exec_command(TradingCommand::SubmitOrder(command));
1414 }
1415 }
1416 }
1417
1418 pub fn fill_market_order(&mut self, client_order_id: ClientOrderId) {
1422 let mut order = match self
1423 .cache
1424 .borrow()
1425 .order(&client_order_id)
1426 .map(|o| o.clone())
1427 {
1428 Some(order) => order,
1429 None => {
1430 log::error!("Cannot fill market order: order {client_order_id} not found in cache");
1431 return;
1432 }
1433 };
1434
1435 let trigger_instrument_id = order
1436 .trigger_instrument_id()
1437 .unwrap_or(order.instrument_id());
1438
1439 let matching_core = match self.matching_cores.get(&trigger_instrument_id) {
1440 Some(core) => core,
1441 None => {
1442 log::error!(
1443 "Cannot fill market order: no matching core for instrument {trigger_instrument_id}"
1444 );
1445 return; }
1447 };
1448
1449 let released_price =
1450 match self.validate_release(&order, matching_core, trigger_instrument_id) {
1451 Some(price) => price,
1452 None => return, };
1454
1455 let command = self
1456 .manager
1457 .pop_submit_order_command(order.client_order_id())
1458 .expect("invalid operation `fill_market_order` with no command");
1459
1460 if let Some(matching_core) = self.matching_cores.get_mut(&trigger_instrument_id) {
1461 if let Err(e) = matching_core.delete_order(client_order_id) {
1462 log::debug!("Cannot delete order: {e:?}");
1463 }
1464
1465 order.set_emulation_trigger(None);
1466
1467 let mut transformed = MarketOrder::new(
1469 order.trader_id(),
1470 order.strategy_id(),
1471 order.instrument_id(),
1472 order.client_order_id(),
1473 order.order_side(),
1474 order.quantity(),
1475 order.time_in_force(),
1476 UUID4::new(),
1477 self.clock.borrow().timestamp_ns(),
1478 order.is_reduce_only(),
1479 order.is_quote_quantity(),
1480 order.contingency_type(),
1481 order.order_list_id(),
1482 order.linked_order_ids().map(Vec::from),
1483 order.parent_order_id(),
1484 order.exec_algorithm_id(),
1485 order.exec_algorithm_params().cloned(),
1486 order.exec_spawn_id(),
1487 order.tags().map(Vec::from),
1488 );
1489
1490 let original_events = order.events();
1491
1492 for event in original_events.into_iter().rev() {
1495 transformed.events.insert(0, event.clone());
1496 }
1497
1498 let add_result = {
1499 let mut cache = self.cache.borrow_mut();
1500 cache.add_order(
1501 OrderAny::Market(transformed.clone()),
1502 command.position_id,
1503 command.client_id,
1504 true,
1505 )
1506 };
1507
1508 if let Err(e) = add_result {
1509 log::error!("Failed to add order: {e}");
1510 } else {
1511 msgbus::publish_order_event(
1512 format!("events.order.{}", order.strategy_id()).into(),
1513 transformed.last_event(),
1514 );
1515 }
1516
1517 let ts_now = self.clock.borrow().timestamp_ns();
1518 let event = OrderReleased::new(
1519 order.trader_id(),
1520 order.strategy_id(),
1521 order.instrument_id(),
1522 order.client_order_id(),
1523 released_price,
1524 UUID4::new(),
1525 ts_now,
1526 ts_now,
1527 );
1528
1529 let event = OrderEventAny::Released(event);
1530
1531 if let Err(e) = self.cache.borrow_mut().update_order(&event) {
1532 log::error!("Failed to apply order event: {e}");
1533 return;
1534 }
1535 self.send_risk_event(event.clone());
1536
1537 log::info!("Releasing order {}", order.client_order_id());
1538
1539 msgbus::publish_order_event(
1541 format!("events.order.{}", order.strategy_id()).into(),
1542 &event,
1543 );
1544
1545 if let Some(exec_algorithm_id) = order.exec_algorithm_id() {
1546 self.send_algo_command(command, exec_algorithm_id);
1547 } else {
1548 self.send_exec_command(TradingCommand::SubmitOrder(command));
1549 }
1550 }
1551 }
1552
1553 fn update_trailing_stop_order(&mut self, order: &mut OrderAny) {
1554 let trigger_instrument_id = order
1555 .trigger_instrument_id()
1556 .unwrap_or_else(|| order.instrument_id());
1557 let Some(matching_core) = self.matching_cores.get(&trigger_instrument_id) else {
1558 log::error!(
1559 "Cannot update trailing-stop order: no matching core for instrument {trigger_instrument_id}"
1560 );
1561 return;
1562 };
1563
1564 let mut bid = matching_core.bid;
1565 let mut ask = matching_core.ask;
1566 let mut last = matching_core.last;
1567 let price_increment = matching_core.price_increment;
1568 let instrument_id = matching_core.instrument_id;
1569
1570 if bid.is_none() || ask.is_none() || last.is_none() {
1571 if let Some(q) = self.cache.borrow().quote(&instrument_id) {
1572 bid.get_or_insert(q.bid_price);
1573 ask.get_or_insert(q.ask_price);
1574 }
1575
1576 if let Some(t) = self.cache.borrow().trade(&instrument_id) {
1577 last.get_or_insert(t.price);
1578 }
1579 }
1580
1581 let was_activated = is_order_activated(order);
1582
1583 if !self.maybe_activate_trailing_stop(order, bid, ask, last) {
1584 return; }
1586
1587 if !was_activated {
1588 self.refresh_matching_core_entry(order, trigger_instrument_id);
1591 }
1592
1593 let (new_trigger_px, new_limit_px) = match trailing_stop_calculate(
1594 price_increment,
1595 order.trigger_price(),
1596 order,
1597 bid,
1598 ask,
1599 last,
1600 ) {
1601 Ok(pair) => pair,
1602 Err(e) => {
1603 log::warn!("Cannot calculate trailing-stop update: {e}");
1604 return;
1605 }
1606 };
1607
1608 if new_trigger_px.is_none() && new_limit_px.is_none() {
1609 return;
1610 }
1611
1612 let ts_now = self.clock.borrow().timestamp_ns();
1613 let update = OrderUpdated::new(
1614 order.trader_id(),
1615 order.strategy_id(),
1616 order.instrument_id(),
1617 order.client_order_id(),
1618 order.quantity(),
1619 UUID4::new(),
1620 ts_now,
1621 ts_now,
1622 false,
1623 order.venue_order_id(),
1624 order.account_id(),
1625 new_limit_px,
1626 new_trigger_px,
1627 None,
1628 order.is_quote_quantity(),
1629 );
1630 let wrapped = OrderEventAny::Updated(update);
1631
1632 *order = match self.cache.borrow_mut().update_order(&wrapped) {
1633 Ok(order) => order,
1634 Err(e) => {
1635 log::error!("Failed to apply order event: {e}");
1636 return;
1637 }
1638 };
1639
1640 self.refresh_matching_core_entry(order, trigger_instrument_id);
1641
1642 self.send_risk_event(wrapped);
1643 }
1644
1645 fn refresh_matching_core_entry(
1646 &mut self,
1647 order: &OrderAny,
1648 trigger_instrument_id: InstrumentId,
1649 ) {
1650 if let Some(matching_core) = self.matching_cores.get_mut(&trigger_instrument_id) {
1651 if let Err(e) = matching_core.delete_order(order.client_order_id()) {
1652 log::debug!("Cannot update trailing-stop match info: {e:?}");
1653 }
1654
1655 let (trigger_price, limit_price) = if order.trigger_price().is_some() {
1658 (order.trigger_price(), order.price())
1659 } else {
1660 (None, None)
1661 };
1662
1663 matching_core.add_order(RestingOrder::new(
1664 order.client_order_id(),
1665 order.order_side(),
1666 order.order_type(),
1667 trigger_price,
1668 limit_price,
1669 is_order_activated(order),
1670 ));
1671 }
1672 }
1673
1674 fn maybe_activate_trailing_stop(
1680 &self,
1681 order: &mut OrderAny,
1682 bid: Option<Price>,
1683 ask: Option<Price>,
1684 last: Option<Price>,
1685 ) -> bool {
1686 let (is_activated, activation_price, trigger_type, order_side) = match order {
1687 OrderAny::TrailingStopMarket(inner) => (
1688 inner.is_activated,
1689 inner.activation_price,
1690 inner.trigger_type,
1691 inner.order_side(),
1692 ),
1693 OrderAny::TrailingStopLimit(inner) => (
1694 inner.is_activated,
1695 inner.activation_price,
1696 inner.trigger_type,
1697 inner.order_side(),
1698 ),
1699 _ => return true,
1700 };
1701
1702 if is_activated {
1703 return true;
1704 }
1705
1706 if let Some(activation_price) = activation_price {
1707 let hit = match order_side {
1708 OrderSide::Buy => ask.is_some_and(|a| a <= activation_price),
1709 OrderSide::Sell => bid.is_some_and(|b| b >= activation_price),
1710 };
1711
1712 if hit {
1713 Self::set_trailing_stop_activated(order, None);
1714 self.persist_trailing_stop_activation(order);
1715 }
1716 return hit;
1717 }
1718
1719 let market_price = match trigger_type {
1720 TriggerType::LastPrice => last,
1721 _ => match order_side {
1722 OrderSide::Buy => ask,
1723 OrderSide::Sell => bid,
1724 },
1725 };
1726
1727 let Some(market_price) = market_price else {
1728 log::error!(
1729 "Cannot activate trailing stop {}: no market price available",
1730 order.client_order_id()
1731 );
1732 return false;
1733 };
1734
1735 Self::set_trailing_stop_activated(order, Some(market_price));
1736 self.persist_trailing_stop_activation(order);
1737 true
1738 }
1739
1740 fn set_trailing_stop_activated(order: &mut OrderAny, activation_price: Option<Price>) {
1741 match order {
1742 OrderAny::TrailingStopMarket(inner) => {
1743 if let Some(price) = activation_price {
1744 inner.activation_price = Some(price);
1745 }
1746 inner.set_activated();
1747 }
1748 OrderAny::TrailingStopLimit(inner) => {
1749 if let Some(price) = activation_price {
1750 inner.activation_price = Some(price);
1751 }
1752 inner.set_activated();
1753 }
1754 _ => {}
1755 }
1756 }
1757
1758 fn persist_trailing_stop_activation(&self, order: &OrderAny) {
1759 if let Err(e) = self.cache.borrow_mut().replace_order(order) {
1760 log::error!("Failed to update order: {e}");
1761 }
1762 }
1763
1764 fn send_algo_command(&self, command: SubmitOrder, exec_algorithm_id: ExecAlgorithmId) {
1765 let id = command.strategy_id;
1766 log::info!("{id} {CMD}{SEND} {command}");
1767
1768 let endpoint = format!("{exec_algorithm_id}.execute");
1769 msgbus::send_any(endpoint.into(), &TradingCommand::SubmitOrder(command));
1770 }
1771
1772 fn send_risk_command(&self, command: TradingCommand) {
1773 log_cmd_send(&command);
1774 let endpoint = MessagingSwitchboard::risk_engine_queue_execute();
1775 msgbus::send_trading_command(endpoint, command);
1776 }
1777
1778 fn send_exec_command(&self, command: TradingCommand) {
1779 log_cmd_send(&command);
1780 let endpoint = MessagingSwitchboard::exec_engine_queue_execute();
1781 msgbus::send_trading_command(endpoint, command);
1782 }
1783
1784 fn send_risk_event(&self, event: OrderEventAny) {
1785 log_evt_send(&event);
1786 let endpoint = MessagingSwitchboard::risk_engine_process();
1787 msgbus::send_order_event(endpoint, event);
1788 }
1789
1790 fn send_exec_event(&self, event: OrderEventAny) {
1791 log_evt_send(&event);
1792 let endpoint = MessagingSwitchboard::exec_engine_process();
1793 msgbus::send_order_event(endpoint, event);
1794 }
1795
1796 fn send_portfolio_order_event(&self, event: OrderEventAny) {
1797 log_evt_send(&event);
1798 let endpoint = MessagingSwitchboard::portfolio_update_order();
1799 msgbus::send_order_event(endpoint, event);
1800 }
1801
1802 fn send_data_command(&self, command: DataCommand) {
1803 log::info!("{CMD}{SEND} {command:?}");
1804 let endpoint = MessagingSwitchboard::data_engine_queue_execute();
1805 msgbus::send_data_command(endpoint, command);
1806 }
1807}
1808
1809fn publish_order_event(event: &OrderEventAny) {
1810 msgbus::publish_order_event(get_event_order_topic(event.strategy_id()), event);
1811
1812 if let OrderEventAny::Canceled(_) = event {
1813 msgbus::publish_order_event(get_order_canceled_topic(event.instrument_id()), event);
1814 }
1815}
1816
1817fn is_order_activated(order: &OrderAny) -> bool {
1818 match order {
1819 OrderAny::TrailingStopMarket(o) => o.is_activated,
1820 OrderAny::TrailingStopLimit(o) => o.is_activated,
1821 _ => true,
1822 }
1823}
1824
1825fn log_cmd_send(command: &TradingCommand) {
1826 if let Some(id) = command.strategy_id() {
1827 log::info!("{id} {CMD}{SEND} {command}");
1828 } else {
1829 log::info!("{CMD}{SEND} {command}");
1830 }
1831}
1832
1833fn log_evt_send(event: &OrderEventAny) {
1834 let id = event.strategy_id();
1835 log::info!("{id} {EVT}{SEND} {event}");
1836}
1837
1838#[cfg(test)]
1839mod tests {
1840 use std::{cell::RefCell, rc::Rc};
1841
1842 use nautilus_common::{
1843 cache::Cache,
1844 clock::TestClock,
1845 messages::data::{DataCommand, SubscribeCommand, UnsubscribeCommand},
1846 msgbus::{
1847 MessagingSwitchboard,
1848 stubs::{
1849 TypedIntoMessageSavingHandler, get_any_saving_handler,
1850 get_typed_into_message_saving_handler,
1851 },
1852 },
1853 };
1854 use nautilus_core::UUID4;
1855 use nautilus_model::{
1856 data::{QuoteTick, TradeTick},
1857 enums::{AggressorSide, OrderSide, OrderType, TrailingOffsetType, TriggerType},
1858 identifiers::{
1859 ClientId, ClientOrderId, OrderListId, StrategyId, Symbol, TradeId, TraderId,
1860 },
1861 instruments::{
1862 CryptoPerpetual, Instrument, InstrumentAny, SyntheticInstrument,
1863 stubs::crypto_perpetual_ethusdt,
1864 },
1865 orders::{OrderList, OrderTestBuilder},
1866 types::{Price, Quantity},
1867 };
1868 use rstest::{fixture, rstest};
1869 use rust_decimal_macros::dec;
1870 use ustr::Ustr;
1871
1872 use super::*;
1873
1874 #[fixture]
1875 fn instrument() -> CryptoPerpetual {
1876 crypto_perpetual_ethusdt()
1877 }
1878
1879 #[expect(clippy::type_complexity)]
1880 fn create_emulator() -> (
1881 Rc<RefCell<dyn Clock>>,
1882 Rc<RefCell<Cache>>,
1883 Rc<RefCell<OrderEmulator>>,
1884 ) {
1885 let clock: Rc<RefCell<dyn Clock>> = Rc::new(RefCell::new(TestClock::new()));
1886 let cache = Rc::new(RefCell::new(Cache::new(None, None)));
1887 let emulator = Rc::new(RefCell::new(OrderEmulator::new(
1888 clock.clone(),
1889 cache.clone(),
1890 )));
1891
1892 OrderEmulator::register_msgbus_handlers(&emulator);
1893
1894 (clock, cache, emulator)
1895 }
1896
1897 fn create_stop_market_order(instrument: &CryptoPerpetual, trigger: TriggerType) -> OrderAny {
1898 OrderTestBuilder::new(OrderType::StopMarket)
1899 .instrument_id(instrument.id())
1900 .side(OrderSide::Buy)
1901 .trigger_price(Price::from("5100.00"))
1902 .quantity(Quantity::from(1))
1903 .emulation_trigger(trigger)
1904 .build()
1905 }
1906
1907 fn create_stop_limit_order(instrument: &CryptoPerpetual, trigger: TriggerType) -> OrderAny {
1908 OrderTestBuilder::new(OrderType::StopLimit)
1909 .instrument_id(instrument.id())
1910 .side(OrderSide::Buy)
1911 .price(Price::from("5100.00"))
1912 .trigger_price(Price::from("5100.00"))
1913 .quantity(Quantity::from(1))
1914 .emulation_trigger(trigger)
1915 .build()
1916 }
1917
1918 fn create_list_stop_market_order(
1919 instrument: &CryptoPerpetual,
1920 client_order_id: &str,
1921 order_list_id: OrderListId,
1922 ) -> OrderAny {
1923 OrderTestBuilder::new(OrderType::StopMarket)
1924 .instrument_id(instrument.id())
1925 .client_order_id(ClientOrderId::from(client_order_id))
1926 .order_list_id(order_list_id)
1927 .side(OrderSide::Buy)
1928 .trigger_price(Price::from("5100.00"))
1929 .quantity(Quantity::from(1))
1930 .emulation_trigger(TriggerType::BidAsk)
1931 .build()
1932 }
1933
1934 fn create_submit_order(instrument: &CryptoPerpetual, order: &OrderAny) -> SubmitOrder {
1935 SubmitOrder::new(
1936 TraderId::from("TRADER-001"),
1937 None,
1938 StrategyId::from("STRATEGY-001"),
1939 instrument.id(),
1940 order.client_order_id(),
1941 order.init_event().clone(),
1942 None,
1943 None,
1944 None,
1945 UUID4::new(),
1946 0.into(),
1947 None, )
1949 }
1950
1951 fn create_quote_tick(instrument: &CryptoPerpetual, bid: &str, ask: &str) -> QuoteTick {
1952 QuoteTick::new(
1953 instrument.id(),
1954 Price::from(bid),
1955 Price::from(ask),
1956 Quantity::from(10),
1957 Quantity::from(10),
1958 0.into(),
1959 0.into(),
1960 )
1961 }
1962
1963 fn create_trade_tick(instrument: &CryptoPerpetual, price: &str) -> TradeTick {
1964 TradeTick::new(
1965 instrument.id(),
1966 Price::from(price),
1967 Quantity::from(1),
1968 AggressorSide::Buy,
1969 TradeId::from("T-001"),
1970 0.into(),
1971 0.into(),
1972 )
1973 }
1974
1975 fn add_instrument_to_cache(cache: &Rc<RefCell<Cache>>, instrument: &CryptoPerpetual) {
1976 cache
1977 .borrow_mut()
1978 .add_instrument(InstrumentAny::CryptoPerpetual(instrument.clone()))
1979 .unwrap();
1980 }
1981
1982 fn register_risk_event_handler(id: &str) -> TypedIntoMessageSavingHandler<OrderEventAny> {
1983 let (handler, saving_handler) =
1984 get_typed_into_message_saving_handler::<OrderEventAny>(Some(Ustr::from(id)));
1985 msgbus::register_order_event_endpoint(MessagingSwitchboard::risk_engine_process(), handler);
1986 saving_handler
1987 }
1988
1989 fn register_exec_event_handler(
1990 cache: Rc<RefCell<Cache>>,
1991 id: &str,
1992 ) -> TypedIntoMessageSavingHandler<OrderEventAny> {
1993 let messages = Rc::new(RefCell::new(Vec::new()));
1994 let messages_for_handler = messages.clone();
1995 msgbus::register_order_event_endpoint(
1996 MessagingSwitchboard::exec_engine_process(),
1997 TypedIntoHandler::from(move |event: OrderEventAny| {
1998 cache.borrow_mut().update_order(&event).unwrap();
1999 messages_for_handler.borrow_mut().push(event);
2000 }),
2001 );
2002 TypedIntoMessageSavingHandler::new_with_messages(Some(Ustr::from(id)), messages)
2003 }
2004
2005 fn register_portfolio_event_handler(id: &str) -> TypedIntoMessageSavingHandler<OrderEventAny> {
2006 let (handler, saving_handler) =
2007 get_typed_into_message_saving_handler::<OrderEventAny>(Some(Ustr::from(id)));
2008 msgbus::register_order_event_endpoint(
2009 MessagingSwitchboard::portfolio_update_order(),
2010 handler,
2011 );
2012 saving_handler
2013 }
2014
2015 fn register_data_command_handler(id: &str) -> TypedIntoMessageSavingHandler<DataCommand> {
2016 let (handler, saving_handler) =
2017 get_typed_into_message_saving_handler::<DataCommand>(Some(Ustr::from(id)));
2018 msgbus::register_data_command_endpoint(
2019 MessagingSwitchboard::data_engine_queue_execute(),
2020 handler,
2021 );
2022 saving_handler
2023 }
2024
2025 fn subscribe_order_topic(
2026 strategy_id: StrategyId,
2027 ) -> (TypedHandler<OrderEventAny>, Rc<RefCell<Vec<OrderEventAny>>>) {
2028 let events = Rc::new(RefCell::new(Vec::new()));
2029 let handler = TypedHandler::from({
2030 let events = events.clone();
2031 move |event: &OrderEventAny| {
2032 events.borrow_mut().push(event.clone());
2033 }
2034 });
2035 msgbus::subscribe_order_events(
2036 format!("events.order.{strategy_id}").into(),
2037 handler.clone(),
2038 None,
2039 );
2040 (handler, events)
2041 }
2042
2043 fn subscribe_order_cancel_topic(
2044 instrument_id: InstrumentId,
2045 ) -> (TypedHandler<OrderEventAny>, Rc<RefCell<Vec<OrderEventAny>>>) {
2046 let events = Rc::new(RefCell::new(Vec::new()));
2047 let handler = TypedHandler::from({
2048 let events = events.clone();
2049 move |event: &OrderEventAny| {
2050 events.borrow_mut().push(event.clone());
2051 }
2052 });
2053 msgbus::subscribe_order_events(
2054 get_order_canceled_topic(instrument_id).into(),
2055 handler.clone(),
2056 None,
2057 );
2058 (handler, events)
2059 }
2060
2061 #[rstest]
2062 fn test_dispatch_manager_publish_initialized_publishes_order_event(
2063 instrument: CryptoPerpetual,
2064 ) {
2065 let (_clock, _cache, emulator) = create_emulator();
2066 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2067 let strategy_id = order.strategy_id();
2068 let client_order_id = order.client_order_id();
2069 let event = OrderEventAny::Initialized(order.init_event().clone());
2070 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
2071
2072 emulator
2073 .borrow_mut()
2074 .dispatch_manager_action(OrderManagerAction::PublishInitialized(event));
2075 msgbus::unsubscribe_order_events(
2076 format!("events.order.{strategy_id}").into(),
2077 &order_handler,
2078 );
2079 let order_events = order_events.borrow();
2080
2081 assert_eq!(order_events.len(), 1);
2082 assert!(matches!(
2083 &order_events[0],
2084 OrderEventAny::Initialized(event) if event.client_order_id == client_order_id
2085 ));
2086 }
2087
2088 #[rstest]
2089 fn test_dispatch_manager_submit_to_emulator_applies_emulated_event(
2090 instrument: CryptoPerpetual,
2091 ) {
2092 let (_clock, cache, emulator) = create_emulator();
2093 let risk_events = register_risk_event_handler("RiskEngine.process.dispatch_emulated");
2094 add_instrument_to_cache(&cache, &instrument);
2095 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2096 let client_order_id = order.client_order_id();
2097 let command = create_submit_order(&instrument, &order);
2098 cache
2099 .borrow_mut()
2100 .add_order(order, None, None, false)
2101 .unwrap();
2102 emulator
2103 .borrow_mut()
2104 .cache_submit_order_command(command.clone());
2105
2106 emulator
2107 .borrow_mut()
2108 .dispatch_manager_action(OrderManagerAction::SubmitToEmulator(command));
2109 let cache = cache.borrow();
2110 let cached_order = cache.order(&client_order_id).unwrap();
2111 let risk_events = risk_events.get_messages();
2112
2113 assert_eq!(cached_order.status(), OrderStatus::Emulated);
2114 assert_eq!(risk_events.len(), 1);
2115 assert!(matches!(risk_events[0], OrderEventAny::Emulated(_)));
2116 }
2117
2118 #[rstest]
2119 fn test_registered_execute_endpoint_routes_submit_order(instrument: CryptoPerpetual) {
2120 let (_clock, cache, emulator) = create_emulator();
2121 let risk_events = register_risk_event_handler("RiskEngine.process.endpoint_emulated");
2122 let data_commands = register_data_command_handler("DataEngine.queue_execute.endpoint");
2123 add_instrument_to_cache(&cache, &instrument);
2124 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2125 let client_order_id = order.client_order_id();
2126 let command = create_submit_order(&instrument, &order);
2127 cache
2128 .borrow_mut()
2129 .add_order(order, None, None, false)
2130 .unwrap();
2131 emulator
2132 .borrow_mut()
2133 .cache_submit_order_command(command.clone());
2134
2135 msgbus::send_trading_command(
2136 MessagingSwitchboard::order_emulator_execute(),
2137 TradingCommand::SubmitOrder(command),
2138 );
2139 let cache = cache.borrow();
2140 let cached_order = cache.order(&client_order_id).unwrap();
2141 let risk_events = risk_events.get_messages();
2142 let data_commands = data_commands.get_messages();
2143
2144 assert_eq!(cached_order.status(), OrderStatus::Emulated);
2145 assert_eq!(risk_events.len(), 1);
2146 assert!(matches!(risk_events[0], OrderEventAny::Emulated(_)));
2147 assert_eq!(data_commands.len(), 1);
2148 assert!(matches!(
2149 &data_commands[0],
2150 DataCommand::Subscribe(SubscribeCommand::Quotes(command))
2151 if command.instrument_id == instrument.id()
2152 ));
2153 }
2154
2155 #[rstest]
2156 fn test_cancel_all_orders_filters_emulated_orders_by_client(instrument: CryptoPerpetual) {
2157 let (_clock, cache, emulator) = create_emulator();
2158 let _risk_events = register_risk_event_handler("RiskEngine.process.cancel_all_client");
2159 let portfolio_events =
2160 register_portfolio_event_handler("Portfolio.update_order.cancel_all_client");
2161 let _data_commands = register_data_command_handler("DataEngine.queue_execute.cancel_all");
2162 add_instrument_to_cache(&cache, &instrument);
2163
2164 let selected_client = ClientId::from("CLIENT-001");
2165 let other_client = ClientId::from("CLIENT-002");
2166 let make_order = |client_order_id| {
2167 OrderTestBuilder::new(OrderType::StopMarket)
2168 .instrument_id(instrument.id())
2169 .client_order_id(client_order_id)
2170 .side(OrderSide::Buy)
2171 .trigger_price(Price::from("5100.00"))
2172 .quantity(Quantity::from(1))
2173 .emulation_trigger(TriggerType::BidAsk)
2174 .build()
2175 };
2176 let selected_order = make_order(ClientOrderId::from("O-EMULATED-SELECTED"));
2177 let other_order = make_order(ClientOrderId::from("O-EMULATED-OTHER"));
2178 let unclaimed_order = make_order(ClientOrderId::from("O-EMULATED-UNCLAIMED"));
2179 let synthetic_formula = format!("{} * 1.0", instrument.id());
2180 let synthetic = SyntheticInstrument::builder()
2181 .symbol(Symbol::from("ETH-INDEX"))
2182 .price_precision(instrument.price_precision())
2183 .components(vec![instrument.id()])
2184 .formula(&synthetic_formula)
2185 .ts_event(0.into())
2186 .ts_init(0.into())
2187 .build()
2188 .unwrap();
2189 let trigger_instrument_id = synthetic.id;
2190 let cross_trigger_order = OrderTestBuilder::new(OrderType::StopMarket)
2191 .instrument_id(instrument.id())
2192 .client_order_id(ClientOrderId::from("O-EMULATED-CROSS-TRIGGER"))
2193 .side(OrderSide::Buy)
2194 .trigger_price(Price::from("5100.00"))
2195 .quantity(Quantity::from(1))
2196 .emulation_trigger(TriggerType::BidAsk)
2197 .trigger_instrument_id(trigger_instrument_id)
2198 .build();
2199 let mut selected_submit = create_submit_order(&instrument, &selected_order);
2200 selected_submit.client_id = Some(selected_client);
2201 let mut other_submit = create_submit_order(&instrument, &other_order);
2202 other_submit.client_id = Some(other_client);
2203 let unclaimed_submit = create_submit_order(&instrument, &unclaimed_order);
2204 let mut cross_trigger_submit = create_submit_order(&instrument, &cross_trigger_order);
2205 cross_trigger_submit.client_id = Some(selected_client);
2206 {
2207 let mut cache = cache.borrow_mut();
2208 cache.add_synthetic(synthetic).unwrap();
2209 cache
2210 .add_order(selected_order.clone(), None, Some(selected_client), false)
2211 .unwrap();
2212 cache
2213 .add_order(other_order.clone(), None, Some(other_client), false)
2214 .unwrap();
2215 cache
2216 .add_order(unclaimed_order.clone(), None, None, false)
2217 .unwrap();
2218 cache
2219 .add_order(
2220 cross_trigger_order.clone(),
2221 None,
2222 Some(selected_client),
2223 false,
2224 )
2225 .unwrap();
2226 }
2227 emulator.borrow_mut().handle_submit_order(&selected_submit);
2228 emulator.borrow_mut().handle_submit_order(&other_submit);
2229 emulator.borrow_mut().handle_submit_order(&unclaimed_submit);
2230 emulator
2231 .borrow_mut()
2232 .handle_submit_order(&cross_trigger_submit);
2233
2234 emulator
2235 .borrow_mut()
2236 .execute(TradingCommand::CancelAllOrders(CancelAllOrders::new(
2237 TraderId::from("TRADER-001"),
2238 Some(selected_client),
2239 StrategyId::from("CALLER-001"),
2240 trigger_instrument_id,
2241 None,
2242 UUID4::new(),
2243 0.into(),
2244 None,
2245 None,
2246 )));
2247
2248 assert_eq!(
2249 cache
2250 .borrow()
2251 .order(&cross_trigger_order.client_order_id())
2252 .unwrap()
2253 .status(),
2254 OrderStatus::Emulated
2255 );
2256
2257 emulator
2258 .borrow_mut()
2259 .execute(TradingCommand::CancelAllOrders(CancelAllOrders::new(
2260 TraderId::from("TRADER-001"),
2261 Some(selected_client),
2262 StrategyId::from("CALLER-001"),
2263 instrument.id(),
2264 None,
2265 UUID4::new(),
2266 0.into(),
2267 None,
2268 None,
2269 )));
2270
2271 let cache = cache.borrow();
2272 assert_eq!(
2273 cache
2274 .order(&selected_order.client_order_id())
2275 .unwrap()
2276 .status(),
2277 OrderStatus::Canceled
2278 );
2279 assert_eq!(
2280 cache
2281 .order(&other_order.client_order_id())
2282 .unwrap()
2283 .status(),
2284 OrderStatus::Emulated
2285 );
2286 assert_eq!(
2287 cache
2288 .order(&unclaimed_order.client_order_id())
2289 .unwrap()
2290 .status(),
2291 OrderStatus::Emulated
2292 );
2293 assert_eq!(
2294 cache
2295 .order(&cross_trigger_order.client_order_id())
2296 .unwrap()
2297 .status(),
2298 OrderStatus::Canceled
2299 );
2300 let portfolio_events = portfolio_events.get_messages();
2301 assert_eq!(portfolio_events.len(), 2);
2302 let canceled_ids: AHashSet<_> = portfolio_events
2303 .iter()
2304 .filter_map(|event| match event {
2305 OrderEventAny::Canceled(event) => Some(event.client_order_id),
2306 _ => None,
2307 })
2308 .collect();
2309 assert_eq!(
2310 canceled_ids,
2311 AHashSet::from_iter([
2312 selected_order.client_order_id(),
2313 cross_trigger_order.client_order_id(),
2314 ])
2315 );
2316 }
2317
2318 #[rstest]
2319 fn test_dispatch_manager_submit_to_risk_uses_risk_queue(instrument: CryptoPerpetual) {
2320 let (_clock, _cache, emulator) = create_emulator();
2321 let (handler, messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
2322 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
2323 msgbus::register_trading_command_endpoint(
2324 MessagingSwitchboard::risk_engine_queue_execute(),
2325 handler,
2326 );
2327 let order = create_stop_market_order(&instrument, TriggerType::Default);
2328 let command = create_submit_order(&instrument, &order);
2329 let client_order_id = command.client_order_id;
2330
2331 emulator
2332 .borrow_mut()
2333 .dispatch_manager_action(OrderManagerAction::SubmitToRisk(command));
2334
2335 let messages = messages.get_messages();
2336 assert_eq!(messages.len(), 1);
2337 assert!(matches!(
2338 messages.first(),
2339 Some(TradingCommand::SubmitOrder(command))
2340 if command.client_order_id == client_order_id
2341 ));
2342 }
2343
2344 #[rstest]
2345 fn test_dispatch_manager_submit_to_algorithm_uses_dynamic_endpoint(
2346 instrument: CryptoPerpetual,
2347 ) {
2348 let (_clock, _cache, emulator) = create_emulator();
2349 let exec_algorithm_id = ExecAlgorithmId::from("ALG-001");
2350 let endpoint = format!("{exec_algorithm_id}.execute");
2351 let (handler, messages) =
2352 get_any_saving_handler::<TradingCommand>(Some(Ustr::from("ALG-001.execute")));
2353 msgbus::register_any(endpoint.into(), handler);
2354 let order = create_stop_market_order(&instrument, TriggerType::Default);
2355 let command = create_submit_order(&instrument, &order);
2356 let client_order_id = command.client_order_id;
2357
2358 emulator
2359 .borrow_mut()
2360 .dispatch_manager_action(OrderManagerAction::SubmitToAlgorithm {
2361 command,
2362 exec_algorithm_id,
2363 });
2364
2365 let messages = messages.get_messages();
2366 assert_eq!(messages.len(), 1);
2367 assert!(matches!(
2368 messages.first(),
2369 Some(TradingCommand::SubmitOrder(command))
2370 if command.client_order_id == client_order_id
2371 ));
2372 }
2373
2374 #[rstest]
2375 fn test_dispatch_manager_cancel_local_applies_and_publishes_event(instrument: CryptoPerpetual) {
2376 let (_clock, cache, emulator) = create_emulator();
2377 let portfolio_events =
2378 register_portfolio_event_handler("Portfolio.update_order.dispatch_canceled");
2379 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2380 let client_order_id = order.client_order_id();
2381 let strategy_id = order.strategy_id();
2382 let instrument_id = order.instrument_id();
2383 let command = create_submit_order(&instrument, &order);
2384 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
2385 let (cancel_handler, cancel_events) = subscribe_order_cancel_topic(instrument_id);
2386 cache
2387 .borrow_mut()
2388 .add_order(order.clone(), None, None, false)
2389 .unwrap();
2390 emulator.borrow_mut().cache_submit_order_command(command);
2391
2392 emulator
2393 .borrow_mut()
2394 .dispatch_manager_action(OrderManagerAction::CancelLocal(order));
2395 msgbus::unsubscribe_order_events(get_event_order_topic(strategy_id).into(), &order_handler);
2396 msgbus::unsubscribe_order_events(
2397 get_order_canceled_topic(instrument_id).into(),
2398 &cancel_handler,
2399 );
2400 let cached_order = cache
2401 .borrow()
2402 .order(&client_order_id)
2403 .map(|order| order.clone())
2404 .unwrap();
2405 let portfolio_events = portfolio_events.get_messages();
2406 let order_events = order_events.borrow();
2407 let cancel_events = cancel_events.borrow();
2408 let commands = emulator.borrow().get_submit_order_commands();
2409
2410 assert_eq!(cached_order.status(), OrderStatus::Canceled);
2411 assert_eq!(portfolio_events.len(), 1);
2412 assert_eq!(order_events.len(), 1);
2413 assert_eq!(cancel_events.len(), 1);
2414 assert!(matches!(
2415 &portfolio_events[0],
2416 OrderEventAny::Canceled(event) if event.client_order_id == client_order_id
2417 ));
2418 assert_eq!(order_events[0].client_order_id(), client_order_id);
2419 assert_eq!(cancel_events[0].client_order_id(), client_order_id);
2420 assert!(!commands.contains_key(&client_order_id));
2421 }
2422
2423 #[rstest]
2424 fn test_dispatch_manager_modify_local_quantity_updates_order(instrument: CryptoPerpetual) {
2425 let (_clock, cache, emulator) = create_emulator();
2426 let risk_events = register_risk_event_handler("RiskEngine.process.dispatch_updated");
2427 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2428 let client_order_id = order.client_order_id();
2429 let quantity = Quantity::from(2);
2430 cache
2431 .borrow_mut()
2432 .add_order(order.clone(), None, None, false)
2433 .unwrap();
2434
2435 emulator
2436 .borrow_mut()
2437 .dispatch_manager_action(OrderManagerAction::ModifyLocalQuantity { order, quantity });
2438 let cache = cache.borrow();
2439 let cached_order = cache.order(&client_order_id).unwrap();
2440 let risk_events = risk_events.get_messages();
2441
2442 assert_eq!(cached_order.quantity(), quantity);
2443 assert_eq!(risk_events.len(), 1);
2444 assert!(matches!(
2445 &risk_events[0],
2446 OrderEventAny::Updated(event) if event.client_order_id == client_order_id
2447 ));
2448 }
2449
2450 #[rstest]
2451 fn test_subscribed_quotes_initially_empty() {
2452 let (_clock, _cache, emulator) = create_emulator();
2453
2454 assert!(emulator.borrow().subscribed_quotes().is_empty());
2455 }
2456
2457 #[rstest]
2458 fn test_subscribed_trades_initially_empty() {
2459 let (_clock, _cache, emulator) = create_emulator();
2460
2461 assert!(emulator.borrow().subscribed_trades().is_empty());
2462 }
2463
2464 #[rstest]
2465 fn test_get_submit_order_commands_initially_empty() {
2466 let (_clock, _cache, emulator) = create_emulator();
2467
2468 assert!(emulator.borrow().get_submit_order_commands().is_empty());
2469 }
2470
2471 #[rstest]
2472 fn test_get_matching_core_returns_none_when_not_created(instrument: CryptoPerpetual) {
2473 let (_clock, _cache, emulator) = create_emulator();
2474
2475 assert!(
2476 emulator
2477 .borrow()
2478 .get_matching_core(&instrument.id())
2479 .is_none()
2480 );
2481 }
2482
2483 #[rstest]
2484 fn test_create_matching_core(instrument: CryptoPerpetual) {
2485 let (_clock, _cache, emulator) = create_emulator();
2486
2487 emulator
2488 .borrow_mut()
2489 .create_matching_core(instrument.id(), instrument.price_increment);
2490
2491 assert!(
2492 emulator
2493 .borrow()
2494 .get_matching_core(&instrument.id())
2495 .is_some()
2496 );
2497 }
2498
2499 #[rstest]
2500 fn test_on_quote_tick_no_matching_core_does_not_panic(instrument: CryptoPerpetual) {
2501 let (_clock, _cache, emulator) = create_emulator();
2502 let quote = create_quote_tick(&instrument, "5060.00", "5070.00");
2503
2504 emulator.borrow_mut().on_quote_tick(quote);
2505 }
2506
2507 #[rstest]
2508 fn test_on_trade_tick_no_matching_core_does_not_panic(instrument: CryptoPerpetual) {
2509 let (_clock, _cache, emulator) = create_emulator();
2510 let trade = create_trade_tick(&instrument, "5065.00");
2511
2512 emulator.borrow_mut().on_trade_tick(trade);
2513 }
2514
2515 #[rstest]
2516 fn test_submit_order_bid_ask_trigger_creates_matching_core(instrument: CryptoPerpetual) {
2517 let (_clock, cache, emulator) = create_emulator();
2518 add_instrument_to_cache(&cache, &instrument);
2519 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2520 let command = create_submit_order(&instrument, &order);
2521 cache
2522 .borrow_mut()
2523 .add_order(order, None, None, false)
2524 .unwrap();
2525
2526 emulator
2527 .borrow_mut()
2528 .cache_submit_order_command(command.clone());
2529 emulator.borrow_mut().handle_submit_order(&command);
2530
2531 assert!(
2532 emulator
2533 .borrow()
2534 .get_matching_core(&instrument.id())
2535 .is_some()
2536 );
2537 }
2538
2539 #[rstest]
2540 fn test_submit_order_bid_ask_trigger_tracks_quote_subscription(instrument: CryptoPerpetual) {
2541 let (_clock, cache, emulator) = create_emulator();
2542 let data_commands = register_data_command_handler("DataEngine.queue_execute.quotes");
2543 add_instrument_to_cache(&cache, &instrument);
2544 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2545 let command = create_submit_order(&instrument, &order);
2546 cache
2547 .borrow_mut()
2548 .add_order(order, None, None, false)
2549 .unwrap();
2550
2551 emulator
2552 .borrow_mut()
2553 .cache_submit_order_command(command.clone());
2554 emulator.borrow_mut().handle_submit_order(&command);
2555
2556 assert_eq!(emulator.borrow().subscribed_quotes(), vec![instrument.id()]);
2557 assert!(emulator.borrow().subscribed_trades().is_empty());
2558
2559 let commands = data_commands.get_messages();
2560 assert_eq!(commands.len(), 1);
2561 assert!(matches!(
2562 &commands[0],
2563 DataCommand::Subscribe(SubscribeCommand::Quotes(command))
2564 if command.instrument_id == instrument.id()
2565 ));
2566
2567 let quote = create_quote_tick(&instrument, "5060.00", "5070.00");
2568 msgbus::publish_quote(get_quotes_topic(instrument.id()), "e);
2569 let core = emulator
2570 .borrow()
2571 .get_matching_core(&instrument.id())
2572 .unwrap();
2573 assert_eq!(core.bid, Some(Price::from("5060.00")));
2574 assert_eq!(core.ask, Some(Price::from("5070.00")));
2575 }
2576
2577 #[rstest]
2578 fn test_submit_order_last_price_trigger_tracks_trade_subscription(instrument: CryptoPerpetual) {
2579 let (_clock, cache, emulator) = create_emulator();
2580 let data_commands = register_data_command_handler("DataEngine.queue_execute.trades");
2581 add_instrument_to_cache(&cache, &instrument);
2582 let order = create_stop_market_order(&instrument, TriggerType::LastPrice);
2583 let command = create_submit_order(&instrument, &order);
2584 cache
2585 .borrow_mut()
2586 .add_order(order, None, None, false)
2587 .unwrap();
2588
2589 emulator
2590 .borrow_mut()
2591 .cache_submit_order_command(command.clone());
2592 emulator.borrow_mut().handle_submit_order(&command);
2593
2594 assert!(emulator.borrow().subscribed_quotes().is_empty());
2595 assert_eq!(emulator.borrow().subscribed_trades(), vec![instrument.id()]);
2596
2597 let commands = data_commands.get_messages();
2598 assert_eq!(commands.len(), 1);
2599 assert!(matches!(
2600 &commands[0],
2601 DataCommand::Subscribe(SubscribeCommand::Trades(command))
2602 if command.instrument_id == instrument.id()
2603 ));
2604
2605 let trade = create_trade_tick(&instrument, "5065.00");
2606 msgbus::publish_trade(get_trades_topic(instrument.id()), &trade);
2607 let core = emulator
2608 .borrow()
2609 .get_matching_core(&instrument.id())
2610 .unwrap();
2611 assert_eq!(core.last, Some(Price::from("5065.00")));
2612 }
2613
2614 #[rstest]
2615 fn test_reset_unsubscribes_market_data_and_clears_state(instrument: CryptoPerpetual) {
2616 let (_clock, cache, emulator) = create_emulator();
2617 let data_commands = register_data_command_handler("DataEngine.queue_execute.reset");
2618 add_instrument_to_cache(&cache, &instrument);
2619 let quote_order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2620 let quote_command = create_submit_order(&instrument, "e_order);
2621 let trade_order = OrderTestBuilder::new(OrderType::StopMarket)
2622 .instrument_id(instrument.id())
2623 .client_order_id(ClientOrderId::from("O-RESET-TRADE"))
2624 .side(OrderSide::Buy)
2625 .trigger_price(Price::from("5100.00"))
2626 .quantity(Quantity::from(1))
2627 .emulation_trigger(TriggerType::LastPrice)
2628 .build();
2629 let trade_command = create_submit_order(&instrument, &trade_order);
2630 cache
2631 .borrow_mut()
2632 .add_order(quote_order, None, None, false)
2633 .unwrap();
2634 cache
2635 .borrow_mut()
2636 .add_order(trade_order, None, None, false)
2637 .unwrap();
2638 emulator
2639 .borrow_mut()
2640 .cache_submit_order_command(quote_command.clone());
2641 emulator.borrow_mut().handle_submit_order("e_command);
2642 emulator
2643 .borrow_mut()
2644 .cache_submit_order_command(trade_command.clone());
2645 emulator.borrow_mut().handle_submit_order(&trade_command);
2646 data_commands.clear();
2647
2648 emulator.borrow_mut().reset();
2649 let commands = data_commands.get_messages();
2650 let emulator_ref = emulator.borrow();
2651
2652 assert!(emulator_ref.subscribed_quotes.is_empty());
2653 assert!(emulator_ref.subscribed_trades.is_empty());
2654 assert!(emulator_ref.subscribed_strategies.is_empty());
2655 assert!(emulator_ref.monitored_positions.is_empty());
2656 assert!(emulator_ref.quote_handlers.is_empty());
2657 assert!(emulator_ref.trade_handlers.is_empty());
2658 assert!(emulator_ref.get_submit_order_commands().is_empty());
2659 assert!(emulator_ref.get_matching_core(&instrument.id()).is_none());
2660 assert!(commands.iter().any(|command| matches!(
2661 command,
2662 DataCommand::Unsubscribe(UnsubscribeCommand::Quotes(command))
2663 if command.instrument_id == instrument.id()
2664 )));
2665 assert!(commands.iter().any(|command| matches!(
2666 command,
2667 DataCommand::Unsubscribe(UnsubscribeCommand::Trades(command))
2668 if command.instrument_id == instrument.id()
2669 )));
2670
2671 drop(emulator_ref);
2672 emulator
2673 .borrow_mut()
2674 .create_matching_core(instrument.id(), instrument.price_increment);
2675 let quote = create_quote_tick(&instrument, "5060.00", "5070.00");
2676 let trade = create_trade_tick(&instrument, "5065.00");
2677 msgbus::publish_quote(get_quotes_topic(instrument.id()), "e);
2678 msgbus::publish_trade(get_trades_topic(instrument.id()), &trade);
2679 let core = emulator
2680 .borrow()
2681 .get_matching_core(&instrument.id())
2682 .unwrap();
2683 assert_eq!(core.bid, None);
2684 assert_eq!(core.ask, None);
2685 assert_eq!(core.last, None);
2686 }
2687
2688 #[rstest]
2689 fn test_reset_unsubscribes_in_sorted_instrument_order(instrument: CryptoPerpetual) {
2690 let (_clock, cache, emulator) = create_emulator();
2691 let data_commands = register_data_command_handler("DataEngine.queue_execute.reset_order");
2692
2693 let ids = [
2696 "SOLUSDT-PERP.BINANCE",
2697 "ADAUSDT-PERP.BINANCE",
2698 "XRPUSDT-PERP.BINANCE",
2699 "BTCUSDT-PERP.BINANCE",
2700 "DOTUSDT-PERP.BINANCE",
2701 ];
2702
2703 for (i, id) in ids.iter().enumerate() {
2704 let mut variant = instrument.clone();
2705 variant.id = InstrumentId::from(*id);
2706 add_instrument_to_cache(&cache, &variant);
2707
2708 let order = OrderTestBuilder::new(OrderType::StopMarket)
2709 .instrument_id(variant.id)
2710 .client_order_id(ClientOrderId::from(format!("O-RESET-{i}").as_str()))
2711 .side(OrderSide::Buy)
2712 .trigger_price(Price::from("5100.00"))
2713 .quantity(Quantity::from(1))
2714 .emulation_trigger(TriggerType::BidAsk)
2715 .build();
2716 let command = create_submit_order(&variant, &order);
2717 cache
2718 .borrow_mut()
2719 .add_order(order, None, None, false)
2720 .unwrap();
2721 emulator
2722 .borrow_mut()
2723 .cache_submit_order_command(command.clone());
2724 emulator.borrow_mut().handle_submit_order(&command);
2725 }
2726 data_commands.clear();
2727
2728 emulator.borrow_mut().reset();
2729
2730 let unsubscribed: Vec<InstrumentId> = data_commands
2731 .get_messages()
2732 .iter()
2733 .filter_map(|command| match command {
2734 DataCommand::Unsubscribe(UnsubscribeCommand::Quotes(command)) => {
2735 Some(command.instrument_id)
2736 }
2737 _ => None,
2738 })
2739 .collect();
2740
2741 let mut expected: Vec<InstrumentId> =
2742 ids.iter().map(|id| InstrumentId::from(*id)).collect();
2743 expected.sort();
2744
2745 assert_eq!(unsubscribed, expected);
2746 }
2747
2748 #[rstest]
2749 fn test_submit_order_list_handles_reentrant_order_events(instrument: CryptoPerpetual) {
2750 let (_clock, cache, emulator) = create_emulator();
2751 let data_commands = register_data_command_handler("DataEngine.queue_execute.list");
2752 add_instrument_to_cache(&cache, &instrument);
2753 let order_list_id = OrderListId::from("OL-EMULATOR-001");
2754 let first_order = create_list_stop_market_order(&instrument, "O-LIST-001", order_list_id);
2755 let second_order = create_list_stop_market_order(&instrument, "O-LIST-002", order_list_id);
2756 let orders = vec![first_order.clone(), second_order.clone()];
2757 let order_list = OrderList::from_orders(&orders, 0.into());
2758 let order_inits = orders
2759 .iter()
2760 .map(|order| order.init_event().clone())
2761 .collect();
2762 cache
2763 .borrow_mut()
2764 .add_order(first_order.clone(), None, None, false)
2765 .unwrap();
2766 cache
2767 .borrow_mut()
2768 .add_order(second_order.clone(), None, None, false)
2769 .unwrap();
2770 let command = SubmitOrderList::new(
2771 TraderId::from("TRADER-001"),
2772 None,
2773 StrategyId::from("STRATEGY-001"),
2774 order_list,
2775 order_inits,
2776 None,
2777 None,
2778 None,
2779 UUID4::new(),
2780 0.into(),
2781 None,
2782 );
2783
2784 emulator
2785 .borrow_mut()
2786 .execute(TradingCommand::SubmitOrderList(command));
2787
2788 let commands = data_commands.get_messages();
2789 let cache = cache.borrow();
2790 let first_status = cache
2791 .order(&first_order.client_order_id())
2792 .unwrap()
2793 .status();
2794 let second_status = cache
2795 .order(&second_order.client_order_id())
2796 .unwrap()
2797 .status();
2798 drop(cache);
2799 let emulator = emulator.borrow();
2800
2801 assert_eq!(first_status, OrderStatus::Emulated);
2802 assert_eq!(second_status, OrderStatus::Emulated);
2803 assert!(emulator.get_matching_core(&instrument.id()).is_some());
2804 assert_eq!(emulator.subscribed_quotes(), vec![instrument.id()]);
2805 assert!(commands.iter().any(|command| matches!(
2806 command,
2807 DataCommand::Subscribe(SubscribeCommand::Quotes(command))
2808 if command.instrument_id == instrument.id()
2809 )));
2810 }
2811
2812 #[rstest]
2813 fn test_submit_order_caches_command(instrument: CryptoPerpetual) {
2814 let (_clock, cache, emulator) = create_emulator();
2815 add_instrument_to_cache(&cache, &instrument);
2816 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2817 let client_order_id = order.client_order_id();
2818 let command = create_submit_order(&instrument, &order);
2819 cache
2820 .borrow_mut()
2821 .add_order(order, None, None, false)
2822 .unwrap();
2823
2824 emulator
2825 .borrow_mut()
2826 .cache_submit_order_command(command.clone());
2827 emulator.borrow_mut().handle_submit_order(&command);
2828
2829 let commands = emulator.borrow().get_submit_order_commands();
2830 assert!(commands.contains_key(&client_order_id));
2831 }
2832
2833 #[rstest]
2834 fn test_handle_submit_order_applies_emulated_event_to_cache(instrument: CryptoPerpetual) {
2835 let (_clock, cache, emulator) = create_emulator();
2836 let risk_events = register_risk_event_handler("RiskEngine.process.emulated");
2837 add_instrument_to_cache(&cache, &instrument);
2838 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2839 let client_order_id = order.client_order_id();
2840 let strategy_id = order.strategy_id();
2841 let command = create_submit_order(&instrument, &order);
2842 cache
2843 .borrow_mut()
2844 .add_order(order, None, None, false)
2845 .unwrap();
2846 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
2847
2848 emulator
2849 .borrow_mut()
2850 .cache_submit_order_command(command.clone());
2851 emulator.borrow_mut().handle_submit_order(&command);
2852 msgbus::unsubscribe_order_events(
2853 format!("events.order.{strategy_id}").into(),
2854 &order_handler,
2855 );
2856 let cache = cache.borrow();
2857 let cached_order = cache.order(&client_order_id).unwrap();
2858 let risk_events = risk_events.get_messages();
2859 let order_events = order_events.borrow();
2860
2861 assert_eq!(cached_order.status(), OrderStatus::Emulated);
2862 assert_eq!(cached_order.event_count(), 2);
2863 assert_eq!(risk_events.len(), 1);
2864 assert!(matches!(risk_events[0], OrderEventAny::Emulated(_)));
2865 assert_eq!(order_events.len(), 1);
2866 assert!(matches!(order_events[0], OrderEventAny::Emulated(_)));
2867 }
2868
2869 #[rstest]
2870 fn test_submit_order_missing_synthetic_trigger_cancels_order(instrument: CryptoPerpetual) {
2871 let (_clock, cache, emulator) = create_emulator();
2872 add_instrument_to_cache(&cache, &instrument);
2873 let synthetic_id = InstrumentId::from("BTC-ETH-INDEX.SYNTH");
2874 let order = OrderTestBuilder::new(OrderType::StopMarket)
2875 .instrument_id(instrument.id())
2876 .side(OrderSide::Buy)
2877 .trigger_price(Price::from("5100.00"))
2878 .quantity(Quantity::from(1))
2879 .emulation_trigger(TriggerType::BidAsk)
2880 .trigger_instrument_id(synthetic_id)
2881 .build();
2882 let client_order_id = order.client_order_id();
2883 let command = create_submit_order(&instrument, &order);
2884 cache
2885 .borrow_mut()
2886 .add_order(order, None, None, false)
2887 .unwrap();
2888
2889 emulator
2890 .borrow_mut()
2891 .cache_submit_order_command(command.clone());
2892 emulator.borrow_mut().handle_submit_order(&command);
2893
2894 let cache = cache.borrow();
2895 let cached_order = cache.order(&client_order_id).unwrap();
2896 let emulator = emulator.borrow();
2897 let commands = emulator.get_submit_order_commands();
2898
2899 assert_eq!(cached_order.status(), OrderStatus::Canceled);
2900 assert!(emulator.get_matching_core(&synthetic_id).is_none());
2901 assert!(emulator.subscribed_quotes().is_empty());
2902 assert!(emulator.subscribed_trades().is_empty());
2903 assert!(!commands.contains_key(&client_order_id));
2904 }
2905
2906 #[rstest]
2907 fn test_submit_order_releases_when_immediately_triggered(instrument: CryptoPerpetual) {
2908 let (_clock, cache, emulator) = create_emulator();
2909 let risk_events =
2910 register_risk_event_handler("RiskEngine.process.submit_immediately_triggered");
2911 add_instrument_to_cache(&cache, &instrument);
2912 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
2913 core.set_bid_raw(Price::from("5099.00"));
2914 core.set_ask_raw(Price::from("5101.00"));
2915 emulator
2916 .borrow_mut()
2917 .matching_cores
2918 .insert(instrument.id(), core);
2919 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2920 let client_order_id = order.client_order_id();
2921 let command = create_submit_order(&instrument, &order);
2922 cache
2923 .borrow_mut()
2924 .add_order(order, None, None, false)
2925 .unwrap();
2926
2927 emulator.borrow_mut().handle_submit_order(&command);
2928
2929 let cache = cache.borrow();
2930 let cached_order = cache.order(&client_order_id).unwrap();
2931 let emulator = emulator.borrow();
2932 let matching_core = emulator.get_matching_core(&instrument.id()).unwrap();
2933 let risk_events = risk_events.get_messages();
2934 assert_eq!(cached_order.status(), OrderStatus::Released);
2935 assert_eq!(cached_order.order_type(), OrderType::Market);
2936 assert!(!matching_core.order_exists(client_order_id));
2937 assert_eq!(emulator.subscribed_strategy_count(), 1);
2938 assert!(
2939 !emulator
2940 .get_submit_order_commands()
2941 .contains_key(&client_order_id)
2942 );
2943 assert_eq!(risk_events.len(), 1);
2944 assert!(matches!(
2945 &risk_events[0],
2946 OrderEventAny::Released(event) if event.client_order_id == client_order_id
2947 ));
2948 }
2949
2950 #[rstest]
2951 fn test_submit_limit_order_releases_when_immediately_fillable(instrument: CryptoPerpetual) {
2952 let (_clock, cache, emulator) = create_emulator();
2953 let risk_events =
2954 register_risk_event_handler("RiskEngine.process.submit_immediately_fillable");
2955 add_instrument_to_cache(&cache, &instrument);
2956 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
2957 core.set_bid_raw(Price::from("5099.00"));
2958 core.set_ask_raw(Price::from("5101.00"));
2959 emulator
2960 .borrow_mut()
2961 .matching_cores
2962 .insert(instrument.id(), core);
2963 let order = OrderTestBuilder::new(OrderType::Limit)
2964 .instrument_id(instrument.id())
2965 .side(OrderSide::Buy)
2966 .price(Price::from("5101.00"))
2967 .quantity(Quantity::from(1))
2968 .emulation_trigger(TriggerType::BidAsk)
2969 .build();
2970 let client_order_id = order.client_order_id();
2971 let command = create_submit_order(&instrument, &order);
2972 cache
2973 .borrow_mut()
2974 .add_order(order, None, None, false)
2975 .unwrap();
2976
2977 emulator.borrow_mut().handle_submit_order(&command);
2978
2979 let cache = cache.borrow();
2980 let cached_order = cache.order(&client_order_id).unwrap();
2981 let emulator = emulator.borrow();
2982 let matching_core = emulator.get_matching_core(&instrument.id()).unwrap();
2983 let risk_events = risk_events.get_messages();
2984 assert_eq!(cached_order.status(), OrderStatus::Released);
2985 assert_eq!(cached_order.order_type(), OrderType::Market);
2986 assert!(!matching_core.order_exists(client_order_id));
2987 assert!(
2988 !emulator
2989 .get_submit_order_commands()
2990 .contains_key(&client_order_id)
2991 );
2992 assert_eq!(risk_events.len(), 1);
2993 assert!(matches!(
2994 &risk_events[0],
2995 OrderEventAny::Released(event) if event.client_order_id == client_order_id
2996 ));
2997 }
2998
2999 #[rstest]
3000 fn test_modify_order_reindexes_matching_state(instrument: CryptoPerpetual) {
3001 let (_clock, cache, emulator) = create_emulator();
3002 add_instrument_to_cache(&cache, &instrument);
3003 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3004 let client_order_id = order.client_order_id();
3005 let command = create_submit_order(&instrument, &order);
3006 cache
3007 .borrow_mut()
3008 .add_order(order.clone(), None, None, false)
3009 .unwrap();
3010 emulator.borrow_mut().handle_submit_order(&command);
3011 let exec_events =
3012 register_exec_event_handler(cache.clone(), "ExecEngine.process.modify_reindex");
3013 let new_trigger = Price::from("5200.00");
3014 let modify = ModifyOrder::new(
3015 order.trader_id(),
3016 None,
3017 order.strategy_id(),
3018 instrument.id(),
3019 client_order_id,
3020 None,
3021 None,
3022 None,
3023 Some(new_trigger),
3024 UUID4::new(),
3025 0.into(),
3026 None,
3027 None,
3028 );
3029
3030 emulator.borrow_mut().handle_modify_order(&modify);
3031
3032 let cache = cache.borrow();
3033 let cached_order = cache.order(&client_order_id).unwrap();
3034 let emulator = emulator.borrow();
3035 let matching_core = emulator.get_matching_core(&instrument.id()).unwrap();
3036 let match_info = matching_core.get_order(client_order_id).unwrap();
3037 let exec_events = exec_events.get_messages();
3038 assert_eq!(cached_order.trigger_price(), Some(new_trigger));
3039 assert_eq!(match_info.trigger_price, Some(new_trigger));
3040 assert_eq!(matching_core.get_orders().len(), 1);
3041 assert_eq!(exec_events.len(), 1);
3042 assert!(matches!(
3043 &exec_events[0],
3044 OrderEventAny::Updated(event) if event.client_order_id == client_order_id
3045 ));
3046 }
3047
3048 #[rstest]
3049 fn test_modify_order_releases_when_updated_trigger_matches(instrument: CryptoPerpetual) {
3050 let (_clock, cache, emulator) = create_emulator();
3051 let risk_events =
3052 register_risk_event_handler("RiskEngine.process.modify_immediately_triggered");
3053 add_instrument_to_cache(&cache, &instrument);
3054 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3055 core.set_bid_raw(Price::from("5099.00"));
3056 core.set_ask_raw(Price::from("5101.00"));
3057 emulator
3058 .borrow_mut()
3059 .matching_cores
3060 .insert(instrument.id(), core);
3061 let order = OrderTestBuilder::new(OrderType::StopMarket)
3062 .instrument_id(instrument.id())
3063 .side(OrderSide::Buy)
3064 .trigger_price(Price::from("5200.00"))
3065 .quantity(Quantity::from(1))
3066 .emulation_trigger(TriggerType::BidAsk)
3067 .build();
3068 let client_order_id = order.client_order_id();
3069 let command = create_submit_order(&instrument, &order);
3070 cache
3071 .borrow_mut()
3072 .add_order(order.clone(), None, None, false)
3073 .unwrap();
3074 emulator.borrow_mut().handle_submit_order(&command);
3075 risk_events.clear();
3076 let exec_events = register_exec_event_handler(
3077 cache.clone(),
3078 "ExecEngine.process.modify_immediately_triggered",
3079 );
3080 let modify = ModifyOrder::new(
3081 order.trader_id(),
3082 None,
3083 order.strategy_id(),
3084 instrument.id(),
3085 client_order_id,
3086 None,
3087 None,
3088 None,
3089 Some(Price::from("5100.00")),
3090 UUID4::new(),
3091 0.into(),
3092 None,
3093 None,
3094 );
3095
3096 emulator.borrow_mut().handle_modify_order(&modify);
3097
3098 let cache = cache.borrow();
3099 let cached_order = cache.order(&client_order_id).unwrap();
3100 let emulator = emulator.borrow();
3101 let matching_core = emulator.get_matching_core(&instrument.id()).unwrap();
3102 let exec_events = exec_events.get_messages();
3103 let risk_events = risk_events.get_messages();
3104 assert_eq!(cached_order.status(), OrderStatus::Released);
3105 assert_eq!(cached_order.order_type(), OrderType::Market);
3106 assert!(!matching_core.order_exists(client_order_id));
3107 assert!(
3108 !emulator
3109 .get_submit_order_commands()
3110 .contains_key(&client_order_id)
3111 );
3112 assert_eq!(exec_events.len(), 1);
3113 assert!(matches!(
3114 &exec_events[0],
3115 OrderEventAny::Updated(event) if event.client_order_id == client_order_id
3116 ));
3117 assert_eq!(risk_events.len(), 1);
3118 assert!(matches!(
3119 &risk_events[0],
3120 OrderEventAny::Released(event) if event.client_order_id == client_order_id
3121 ));
3122 }
3123
3124 #[rstest]
3125 fn test_update_order_applies_updated_event_to_cache(instrument: CryptoPerpetual) {
3126 let (_clock, cache, emulator) = create_emulator();
3127 let risk_events = register_risk_event_handler("RiskEngine.process.updated");
3128 let mut order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3129 let client_order_id = order.client_order_id();
3130 cache
3131 .borrow_mut()
3132 .add_order(order.clone(), None, None, false)
3133 .unwrap();
3134
3135 emulator
3136 .borrow_mut()
3137 .update_order(&mut order, Quantity::from(2));
3138 let cache = cache.borrow();
3139 let cached_order = cache.order(&client_order_id).unwrap();
3140 let risk_events = risk_events.get_messages();
3141
3142 assert_eq!(order.quantity(), Quantity::from(2));
3143 assert_eq!(cached_order.quantity(), Quantity::from(2));
3144 assert_eq!(cached_order.status(), OrderStatus::Initialized);
3145 assert_eq!(risk_events.len(), 1);
3146 assert!(matches!(risk_events[0], OrderEventAny::Updated(_)));
3147 }
3148
3149 #[rstest]
3150 fn test_fill_market_order_applies_released_event_to_cache(instrument: CryptoPerpetual) {
3151 let (_clock, cache, emulator) = create_emulator();
3152 let risk_events = register_risk_event_handler("RiskEngine.process.released");
3153 add_instrument_to_cache(&cache, &instrument);
3154 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3155 let client_order_id = order.client_order_id();
3156 let strategy_id = order.strategy_id();
3157 let command = create_submit_order(&instrument, &order);
3158 cache
3159 .borrow_mut()
3160 .add_order(order, None, None, false)
3161 .unwrap();
3162
3163 emulator
3164 .borrow_mut()
3165 .cache_submit_order_command(command.clone());
3166 emulator.borrow_mut().handle_submit_order(&command);
3167 risk_events.clear();
3168 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
3169 {
3170 let mut emulator = emulator.borrow_mut();
3171 emulator
3172 .matching_cores
3173 .get_mut(&instrument.id())
3174 .unwrap()
3175 .set_ask_raw(Price::from("5100.00"));
3176 emulator.fill_market_order(client_order_id);
3177 }
3178 msgbus::unsubscribe_order_events(
3179 format!("events.order.{strategy_id}").into(),
3180 &order_handler,
3181 );
3182 let cache = cache.borrow();
3183 let cached_order = cache.order(&client_order_id).unwrap();
3184 let risk_events = risk_events.get_messages();
3185 let order_events = order_events.borrow();
3186
3187 assert_eq!(cached_order.status(), OrderStatus::Released);
3188 assert_eq!(risk_events.len(), 1);
3189 assert!(matches!(risk_events[0], OrderEventAny::Released(_)));
3190 assert_eq!(order_events.len(), 2);
3191 assert!(matches!(order_events[0], OrderEventAny::Initialized(_)));
3192 assert!(matches!(order_events[1], OrderEventAny::Released(_)));
3193 }
3194
3195 #[rstest]
3196 fn test_fill_limit_order_publishes_transformed_initialized_before_released(
3197 instrument: CryptoPerpetual,
3198 ) {
3199 let (_clock, cache, emulator) = create_emulator();
3200 let risk_events = register_risk_event_handler("RiskEngine.process.limit_released");
3201 add_instrument_to_cache(&cache, &instrument);
3202 let order = create_stop_limit_order(&instrument, TriggerType::BidAsk);
3203 let client_order_id = order.client_order_id();
3204 let strategy_id = order.strategy_id();
3205 let command = create_submit_order(&instrument, &order);
3206 cache
3207 .borrow_mut()
3208 .add_order(order, None, None, false)
3209 .unwrap();
3210
3211 emulator
3212 .borrow_mut()
3213 .cache_submit_order_command(command.clone());
3214 emulator.borrow_mut().handle_submit_order(&command);
3215 risk_events.clear();
3216 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
3217 {
3218 let mut emulator = emulator.borrow_mut();
3219 emulator
3220 .matching_cores
3221 .get_mut(&instrument.id())
3222 .unwrap()
3223 .set_ask_raw(Price::from("5100.00"));
3224 emulator.fill_limit_order(client_order_id);
3225 }
3226 msgbus::unsubscribe_order_events(
3227 format!("events.order.{strategy_id}").into(),
3228 &order_handler,
3229 );
3230 let cache = cache.borrow();
3231 let cached_order = cache.order(&client_order_id).unwrap();
3232 let risk_events = risk_events.get_messages();
3233 let order_events = order_events.borrow();
3234
3235 assert_eq!(cached_order.status(), OrderStatus::Released);
3236 assert_eq!(risk_events.len(), 1);
3237 assert!(matches!(risk_events[0], OrderEventAny::Released(_)));
3238 assert_eq!(order_events.len(), 2);
3239 assert!(matches!(order_events[0], OrderEventAny::Initialized(_)));
3240 assert!(matches!(order_events[1], OrderEventAny::Released(_)));
3241 }
3242
3243 #[rstest]
3244 fn test_quote_tick_updates_matching_core_prices(instrument: CryptoPerpetual) {
3245 let (_clock, cache, emulator) = create_emulator();
3246 add_instrument_to_cache(&cache, &instrument);
3247 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3248 let command = create_submit_order(&instrument, &order);
3249 cache
3250 .borrow_mut()
3251 .add_order(order, None, None, false)
3252 .unwrap();
3253 emulator
3254 .borrow_mut()
3255 .cache_submit_order_command(command.clone());
3256 emulator.borrow_mut().handle_submit_order(&command);
3257
3258 let quote = create_quote_tick(&instrument, "5060.00", "5070.00");
3259 emulator.borrow_mut().on_quote_tick(quote);
3260
3261 let core = emulator
3262 .borrow()
3263 .get_matching_core(&instrument.id())
3264 .unwrap();
3265 assert_eq!(core.bid, Some(Price::from("5060.00")));
3266 assert_eq!(core.ask, Some(Price::from("5070.00")));
3267 }
3268
3269 #[rstest]
3270 fn test_trade_tick_updates_matching_core_last_price(instrument: CryptoPerpetual) {
3271 let (_clock, cache, emulator) = create_emulator();
3272 add_instrument_to_cache(&cache, &instrument);
3273 let order = create_stop_market_order(&instrument, TriggerType::LastPrice);
3274 let command = create_submit_order(&instrument, &order);
3275 cache
3276 .borrow_mut()
3277 .add_order(order, None, None, false)
3278 .unwrap();
3279 emulator
3280 .borrow_mut()
3281 .cache_submit_order_command(command.clone());
3282 emulator.borrow_mut().handle_submit_order(&command);
3283
3284 let trade = create_trade_tick(&instrument, "5065.00");
3285 emulator.borrow_mut().on_trade_tick(trade);
3286
3287 let core = emulator
3288 .borrow()
3289 .get_matching_core(&instrument.id())
3290 .unwrap();
3291 assert_eq!(core.last, Some(Price::from("5065.00")));
3292 }
3293
3294 #[rstest]
3295 fn test_cancel_order_removes_from_matching_core(instrument: CryptoPerpetual) {
3296 let (_clock, cache, emulator) = create_emulator();
3297 add_instrument_to_cache(&cache, &instrument);
3298 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3299 let command = create_submit_order(&instrument, &order);
3300 cache
3301 .borrow_mut()
3302 .add_order(order.clone(), None, None, false)
3303 .unwrap();
3304 emulator
3305 .borrow_mut()
3306 .cache_submit_order_command(command.clone());
3307 emulator.borrow_mut().handle_submit_order(&command);
3308
3309 emulator.borrow_mut().cancel_order(&order);
3310
3311 let core = emulator
3312 .borrow()
3313 .get_matching_core(&instrument.id())
3314 .unwrap();
3315 assert!(core.get_orders().is_empty());
3316 }
3317
3318 #[rstest]
3319 fn test_cancel_order_removes_cached_command(instrument: CryptoPerpetual) {
3320 let (_clock, cache, emulator) = create_emulator();
3321 add_instrument_to_cache(&cache, &instrument);
3322 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3323 let client_order_id = order.client_order_id();
3324 let command = create_submit_order(&instrument, &order);
3325 cache
3326 .borrow_mut()
3327 .add_order(order.clone(), None, None, false)
3328 .unwrap();
3329 emulator
3330 .borrow_mut()
3331 .cache_submit_order_command(command.clone());
3332 emulator.borrow_mut().handle_submit_order(&command);
3333
3334 emulator.borrow_mut().cancel_order(&order);
3335
3336 let commands = emulator.borrow().get_submit_order_commands();
3337 assert!(!commands.contains_key(&client_order_id));
3338 }
3339
3340 #[rstest]
3341 fn test_trailing_stop_waits_for_activation_price(instrument: CryptoPerpetual) {
3342 let (_clock, cache, emulator) = create_emulator();
3343 let risk_events = register_risk_event_handler("RiskEngine.process.trailing_activation");
3344 let (handler, exec_commands): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3345 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3346 msgbus::register_trading_command_endpoint(
3347 MessagingSwitchboard::exec_engine_queue_execute(),
3348 handler,
3349 );
3350 add_instrument_to_cache(&cache, &instrument);
3351
3352 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3353 core.set_bid_raw(Price::from("5000.00"));
3354 core.set_ask_raw(Price::from("5001.00"));
3355 emulator
3356 .borrow_mut()
3357 .matching_cores
3358 .insert(instrument.id(), core);
3359
3360 let order = OrderTestBuilder::new(OrderType::TrailingStopMarket)
3361 .instrument_id(instrument.id())
3362 .client_order_id(ClientOrderId::from("O-TRAILING-1"))
3363 .side(OrderSide::Sell)
3364 .quantity(Quantity::from(1))
3365 .activation_price(Price::from("5050.00"))
3366 .trailing_offset(dec!(10))
3367 .trailing_offset_type(TrailingOffsetType::Price)
3368 .emulation_trigger(TriggerType::BidAsk)
3369 .build();
3370 let client_order_id = order.client_order_id();
3371 let command = create_submit_order(&instrument, &order);
3372 cache
3373 .borrow_mut()
3374 .add_order(order, None, None, false)
3375 .unwrap();
3376
3377 emulator
3378 .borrow_mut()
3379 .cache_submit_order_command(command.clone());
3380 emulator.borrow_mut().handle_submit_order(&command);
3381
3382 {
3383 let cache = cache.borrow();
3384 let cached = cache.order(&client_order_id).unwrap();
3385 let is_activated = match &*cached {
3386 OrderAny::TrailingStopMarket(o) => o.is_activated,
3387 _ => panic!("expected trailing stop market order"),
3388 };
3389 assert!(
3390 !is_activated,
3391 "must not activate before the activation price trades"
3392 );
3393 assert!(cached.trigger_price().is_none());
3394 }
3395 assert!(
3396 !risk_events
3397 .get_messages()
3398 .iter()
3399 .any(|e| matches!(e, OrderEventAny::Updated(_))),
3400 "no trailing update before activation"
3401 );
3402
3403 emulator
3405 .borrow_mut()
3406 .on_quote_tick(create_quote_tick(&instrument, "5055.00", "5056.00"));
3407
3408 {
3409 let cache = cache.borrow();
3410 let cached = cache.order(&client_order_id).unwrap();
3411 let is_activated = match &*cached {
3412 OrderAny::TrailingStopMarket(o) => o.is_activated,
3413 _ => panic!("expected trailing stop market order"),
3414 };
3415 assert!(is_activated);
3416 assert_eq!(cached.trigger_price(), Some(Price::from("5045.00")));
3417 }
3418 assert!(exec_commands.get_messages().is_empty());
3419
3420 emulator
3422 .borrow_mut()
3423 .on_quote_tick(create_quote_tick(&instrument, "5044.00", "5045.00"));
3424
3425 let commands = exec_commands.get_messages();
3426 assert_eq!(commands.len(), 1);
3427 assert!(matches!(
3428 &commands[0],
3429 TradingCommand::SubmitOrder(command) if command.client_order_id == client_order_id
3430 ));
3431 }
3432
3433 #[rstest]
3434 fn test_pending_activation_trailing_stop_limit_not_matched_as_plain_limit(
3435 instrument: CryptoPerpetual,
3436 ) {
3437 let (_clock, cache, emulator) = create_emulator();
3438 let (handler, exec_commands): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3439 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3440 msgbus::register_trading_command_endpoint(
3441 MessagingSwitchboard::exec_engine_queue_execute(),
3442 handler,
3443 );
3444 add_instrument_to_cache(&cache, &instrument);
3445
3446 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3448 core.set_bid_raw(Price::from("5000.00"));
3449 core.set_ask_raw(Price::from("5001.00"));
3450 emulator
3451 .borrow_mut()
3452 .matching_cores
3453 .insert(instrument.id(), core);
3454
3455 let order = OrderTestBuilder::new(OrderType::TrailingStopLimit)
3456 .instrument_id(instrument.id())
3457 .client_order_id(ClientOrderId::from("O-TRAILING-LIMIT-1"))
3458 .side(OrderSide::Sell)
3459 .price(Price::from("5010.00"))
3460 .quantity(Quantity::from(1))
3461 .activation_price(Price::from("5050.00"))
3462 .trailing_offset(dec!(10))
3463 .trailing_offset_type(TrailingOffsetType::Price)
3464 .limit_offset(dec!(5))
3465 .emulation_trigger(TriggerType::BidAsk)
3466 .build();
3467 let client_order_id = order.client_order_id();
3468 let command = create_submit_order(&instrument, &order);
3469 cache
3470 .borrow_mut()
3471 .add_order(order, None, None, false)
3472 .unwrap();
3473
3474 emulator
3475 .borrow_mut()
3476 .cache_submit_order_command(command.clone());
3477 emulator.borrow_mut().handle_submit_order(&command);
3478
3479 emulator
3482 .borrow_mut()
3483 .on_quote_tick(create_quote_tick(&instrument, "4985.00", "4986.00"));
3484
3485 assert!(exec_commands.get_messages().is_empty());
3486 assert!(
3487 emulator
3488 .borrow()
3489 .get_submit_order_commands()
3490 .contains_key(&client_order_id),
3491 "order must remain held by the emulator before activation"
3492 );
3493 }
3494
3495 #[rstest]
3496 fn test_trailing_stop_with_preset_trigger_activates_and_triggers(instrument: CryptoPerpetual) {
3497 let (_clock, cache, emulator) = create_emulator();
3498 let (handler, exec_commands): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3499 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3500 msgbus::register_trading_command_endpoint(
3501 MessagingSwitchboard::exec_engine_queue_execute(),
3502 handler,
3503 );
3504 add_instrument_to_cache(&cache, &instrument);
3505
3506 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3507 core.set_bid_raw(Price::from("5000.00"));
3508 core.set_ask_raw(Price::from("5001.00"));
3509 emulator
3510 .borrow_mut()
3511 .matching_cores
3512 .insert(instrument.id(), core);
3513
3514 let order = OrderTestBuilder::new(OrderType::TrailingStopMarket)
3516 .instrument_id(instrument.id())
3517 .client_order_id(ClientOrderId::from("O-TRAILING-2"))
3518 .side(OrderSide::Sell)
3519 .quantity(Quantity::from(1))
3520 .trigger_price(Price::from("5045.00"))
3521 .activation_price(Price::from("5050.00"))
3522 .trailing_offset(dec!(10))
3523 .trailing_offset_type(TrailingOffsetType::Price)
3524 .emulation_trigger(TriggerType::BidAsk)
3525 .build();
3526 let client_order_id = order.client_order_id();
3527 let command = create_submit_order(&instrument, &order);
3528 cache
3529 .borrow_mut()
3530 .add_order(order, None, None, false)
3531 .unwrap();
3532
3533 emulator
3534 .borrow_mut()
3535 .cache_submit_order_command(command.clone());
3536 emulator.borrow_mut().handle_submit_order(&command);
3537
3538 assert!(
3539 exec_commands.get_messages().is_empty(),
3540 "preset trigger must not release before activation"
3541 );
3542
3543 emulator
3546 .borrow_mut()
3547 .on_quote_tick(create_quote_tick(&instrument, "5055.00", "5056.00"));
3548
3549 assert!(exec_commands.get_messages().is_empty());
3550
3551 emulator
3553 .borrow_mut()
3554 .on_quote_tick(create_quote_tick(&instrument, "5044.00", "5045.00"));
3555
3556 let commands = exec_commands.get_messages();
3557 assert_eq!(commands.len(), 1);
3558 assert!(matches!(
3559 &commands[0],
3560 TradingCommand::SubmitOrder(command) if command.client_order_id == client_order_id
3561 ));
3562 }
3563
3564 #[rstest]
3565 fn test_trailing_stop_activates_despite_trailing_calculate_error(instrument: CryptoPerpetual) {
3566 let (_clock, cache, emulator) = create_emulator();
3567 let (handler, exec_commands): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3568 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3569 msgbus::register_trading_command_endpoint(
3570 MessagingSwitchboard::exec_engine_queue_execute(),
3571 handler,
3572 );
3573 add_instrument_to_cache(&cache, &instrument);
3574
3575 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3577 core.set_bid_raw(Price::from("5000.00"));
3578 core.set_ask_raw(Price::from("5001.00"));
3579 emulator
3580 .borrow_mut()
3581 .matching_cores
3582 .insert(instrument.id(), core);
3583
3584 let order = OrderTestBuilder::new(OrderType::TrailingStopMarket)
3585 .instrument_id(instrument.id())
3586 .client_order_id(ClientOrderId::from("O-TRAILING-3"))
3587 .side(OrderSide::Sell)
3588 .quantity(Quantity::from(1))
3589 .trigger_price(Price::from("5045.00"))
3590 .trigger_type(TriggerType::LastOrBidAsk)
3591 .activation_price(Price::from("5050.00"))
3592 .trailing_offset(dec!(10))
3593 .trailing_offset_type(TrailingOffsetType::Price)
3594 .emulation_trigger(TriggerType::BidAsk)
3595 .build();
3596 let client_order_id = order.client_order_id();
3597 let command = create_submit_order(&instrument, &order);
3598 cache
3599 .borrow_mut()
3600 .add_order(order, None, None, false)
3601 .unwrap();
3602
3603 emulator
3604 .borrow_mut()
3605 .cache_submit_order_command(command.clone());
3606 emulator.borrow_mut().handle_submit_order(&command);
3607
3608 emulator
3611 .borrow_mut()
3612 .on_quote_tick(create_quote_tick(&instrument, "5055.00", "5056.00"));
3613 emulator
3614 .borrow_mut()
3615 .on_quote_tick(create_quote_tick(&instrument, "5044.00", "5045.00"));
3616
3617 let commands = exec_commands.get_messages();
3618 assert_eq!(commands.len(), 1);
3619 assert!(matches!(
3620 &commands[0],
3621 TradingCommand::SubmitOrder(command) if command.client_order_id == client_order_id
3622 ));
3623 }
3624
3625 #[rstest]
3626 fn test_cancel_emulated_oco_leg_cancels_sibling(instrument: CryptoPerpetual) {
3627 let (_clock, cache, emulator) = create_emulator();
3628 add_instrument_to_cache(&cache, &instrument);
3629 let client_order_id_a = ClientOrderId::from("O-OCO-A");
3630 let client_order_id_b = ClientOrderId::from("O-OCO-B");
3631 let order_a = OrderTestBuilder::new(OrderType::StopMarket)
3632 .instrument_id(instrument.id())
3633 .strategy_id(StrategyId::from("STRATEGY-001"))
3634 .client_order_id(client_order_id_a)
3635 .side(OrderSide::Sell)
3636 .trigger_price(Price::from("4900.00"))
3637 .quantity(Quantity::from(1))
3638 .emulation_trigger(TriggerType::BidAsk)
3639 .contingency_type(ContingencyType::Oco)
3640 .linked_order_ids(vec![client_order_id_b])
3641 .build();
3642 let order_b = OrderTestBuilder::new(OrderType::StopMarket)
3643 .instrument_id(instrument.id())
3644 .strategy_id(StrategyId::from("STRATEGY-001"))
3645 .client_order_id(client_order_id_b)
3646 .side(OrderSide::Sell)
3647 .trigger_price(Price::from("4950.00"))
3648 .quantity(Quantity::from(1))
3649 .emulation_trigger(TriggerType::BidAsk)
3650 .contingency_type(ContingencyType::Oco)
3651 .linked_order_ids(vec![client_order_id_a])
3652 .build();
3653 cache
3654 .borrow_mut()
3655 .add_order(order_a.clone(), None, None, false)
3656 .unwrap();
3657 cache
3658 .borrow_mut()
3659 .add_order(order_b.clone(), None, None, false)
3660 .unwrap();
3661
3662 for order in [&order_a, &order_b] {
3663 let command = create_submit_order(&instrument, order);
3664 emulator
3665 .borrow_mut()
3666 .cache_submit_order_command(command.clone());
3667 emulator.borrow_mut().handle_submit_order(&command);
3668 }
3669 assert_eq!(
3670 cache.borrow().order(&client_order_id_a).unwrap().status(),
3671 OrderStatus::Emulated
3672 );
3673 assert_eq!(
3674 cache.borrow().order(&client_order_id_b).unwrap().status(),
3675 OrderStatus::Emulated
3676 );
3677
3678 let cancel = CancelOrder::new(
3679 order_a.trader_id(),
3680 None,
3681 order_a.strategy_id(),
3682 instrument.id(),
3683 client_order_id_a,
3684 None,
3685 UUID4::new(),
3686 0.into(),
3687 None,
3688 None,
3689 );
3690 msgbus::send_trading_command(
3691 MessagingSwitchboard::order_emulator_execute(),
3692 TradingCommand::CancelOrder(cancel),
3693 );
3694
3695 let cache = cache.borrow();
3696 assert_eq!(
3697 cache.order(&client_order_id_a).unwrap().status(),
3698 OrderStatus::Canceled
3699 );
3700 assert_eq!(
3701 cache.order(&client_order_id_b).unwrap().status(),
3702 OrderStatus::Canceled,
3703 "OCO sibling must be canceled through the deferred self-published event"
3704 );
3705 }
3706
3707 #[rstest]
3708 fn test_reentrant_command_from_event_handler_does_not_panic(instrument: CryptoPerpetual) {
3709 let (_clock, cache, emulator) = create_emulator();
3710 add_instrument_to_cache(&cache, &instrument);
3711 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3712 let client_order_id = order.client_order_id();
3713 let strategy_id = order.strategy_id();
3714 let command = create_submit_order(&instrument, &order);
3715 cache
3716 .borrow_mut()
3717 .add_order(order.clone(), None, None, false)
3718 .unwrap();
3719 emulator
3720 .borrow_mut()
3721 .cache_submit_order_command(command.clone());
3722
3723 let cancel = CancelOrder::new(
3726 order.trader_id(),
3727 None,
3728 strategy_id,
3729 instrument.id(),
3730 client_order_id,
3731 None,
3732 UUID4::new(),
3733 0.into(),
3734 None,
3735 None,
3736 );
3737 let handler = TypedHandler::from(move |event: &OrderEventAny| {
3738 if matches!(event, OrderEventAny::Emulated(_)) {
3739 msgbus::send_trading_command(
3740 MessagingSwitchboard::order_emulator_execute(),
3741 TradingCommand::CancelOrder(cancel.clone()),
3742 );
3743 }
3744 });
3745 msgbus::subscribe_order_events(format!("events.order.{strategy_id}").into(), handler, None);
3746
3747 msgbus::send_trading_command(
3749 MessagingSwitchboard::order_emulator_execute(),
3750 TradingCommand::SubmitOrder(command),
3751 );
3752
3753 assert!(
3754 cache.borrow().order(&client_order_id).unwrap().is_closed(),
3755 "deferred cancel must be processed after the active call completes"
3756 );
3757 }
3758
3759 #[rstest]
3760 fn test_released_order_preserves_chronological_event_history(instrument: CryptoPerpetual) {
3761 let (_clock, cache, emulator) = create_emulator();
3762 add_instrument_to_cache(&cache, &instrument);
3763 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3764 core.set_bid_raw(Price::from("5099.00"));
3765 core.set_ask_raw(Price::from("5101.00"));
3766 emulator
3767 .borrow_mut()
3768 .matching_cores
3769 .insert(instrument.id(), core);
3770
3771 let order = create_stop_limit_order(&instrument, TriggerType::BidAsk);
3772 let client_order_id = order.client_order_id();
3773 let original_init = OrderEventAny::Initialized(order.init_event().clone());
3774 let command = create_submit_order(&instrument, &order);
3775 cache
3776 .borrow_mut()
3777 .add_order(order, None, None, false)
3778 .unwrap();
3779
3780 emulator
3781 .borrow_mut()
3782 .cache_submit_order_command(command.clone());
3783 emulator.borrow_mut().handle_submit_order(&command);
3784
3785 let cache = cache.borrow();
3786 let transformed = cache.order(&client_order_id).unwrap();
3787 let events = transformed.events();
3788
3789 assert_eq!(
3790 events.first(),
3791 Some(&&original_init),
3792 "released order history must start with the original OrderInitialized event, was {events:?}",
3793 );
3794 }
3795
3796 #[rstest]
3797 fn test_on_start_cancels_emulated_child_of_closed_positionless_parent(
3798 instrument: CryptoPerpetual,
3799 ) {
3800 let (_clock, cache, emulator) = create_emulator();
3801 add_instrument_to_cache(&cache, &instrument);
3802
3803 let mut parent = OrderTestBuilder::new(OrderType::Limit)
3804 .instrument_id(instrument.id())
3805 .client_order_id(ClientOrderId::from("O-PARENT"))
3806 .side(OrderSide::Buy)
3807 .price(Price::from("5000.00"))
3808 .quantity(Quantity::from(1))
3809 .submit(true)
3810 .build();
3811 parent
3812 .apply(OrderEventAny::Canceled(OrderCanceled::new(
3813 parent.trader_id(),
3814 parent.strategy_id(),
3815 parent.instrument_id(),
3816 parent.client_order_id(),
3817 UUID4::new(),
3818 0.into(),
3819 0.into(),
3820 false,
3821 None,
3822 None,
3823 )))
3824 .unwrap();
3825 assert!(parent.is_closed());
3826 assert!(parent.position_id().is_none());
3827
3828 let mut child = OrderTestBuilder::new(OrderType::StopMarket)
3829 .instrument_id(instrument.id())
3830 .client_order_id(ClientOrderId::from("O-CHILD"))
3831 .side(OrderSide::Sell)
3832 .trigger_price(Price::from("4900.00"))
3833 .quantity(Quantity::from(1))
3834 .emulation_trigger(TriggerType::BidAsk)
3835 .parent_order_id(parent.client_order_id())
3836 .build();
3837 child
3838 .apply(OrderEventAny::Emulated(OrderEmulated::new(
3839 child.trader_id(),
3840 child.strategy_id(),
3841 child.instrument_id(),
3842 child.client_order_id(),
3843 UUID4::new(),
3844 0.into(),
3845 0.into(),
3846 )))
3847 .unwrap();
3848 assert_eq!(child.status(), OrderStatus::Emulated);
3849
3850 cache
3851 .borrow_mut()
3852 .add_order(parent.clone(), None, None, false)
3853 .unwrap();
3854 cache
3855 .borrow_mut()
3856 .add_order(child.clone(), None, None, false)
3857 .unwrap();
3858
3859 emulator.borrow_mut().on_start().unwrap();
3860
3861 assert!(
3862 cache
3863 .borrow()
3864 .order(&child.client_order_id())
3865 .unwrap()
3866 .is_closed(),
3867 "emulated child of a closed position-less parent must be canceled on reactivation"
3868 );
3869 }
3870}