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 None,
1164 );
1165
1166 let event = OrderEventAny::Canceled(event);
1167 if let Err(e) = self.cache.borrow_mut().update_order(&event) {
1168 log::error!("Failed to apply order event: {e}");
1169 return;
1170 }
1171
1172 self.send_portfolio_order_event(event.clone());
1173 publish_order_event(&event);
1174 }
1175
1176 fn check_monitoring(&mut self, strategy_id: StrategyId, position_id: Option<PositionId>) {
1177 if !self.subscribed_strategies.contains(&strategy_id) {
1178 if let Some(handler) = &self.on_event_handler {
1180 msgbus::subscribe_order_events(
1181 format!("events.order.{strategy_id}").into(),
1182 handler.clone(),
1183 None,
1184 );
1185 self.subscribed_strategies.insert(strategy_id);
1186 log::info!("Subscribed to strategy {strategy_id} order events");
1187 }
1188 }
1189
1190 if let Some(position_id) = position_id
1191 && !self.monitored_positions.contains(&position_id)
1192 {
1193 self.monitored_positions.insert(position_id);
1194 }
1195 }
1196
1197 fn validate_release(
1205 &self,
1206 order: &OrderAny,
1207 matching_core: &OrderMatchingCore,
1208 trigger_instrument_id: InstrumentId,
1209 ) -> Option<Price> {
1210 let released_price = match order.order_side() {
1211 OrderSide::Buy => matching_core.ask,
1212 OrderSide::Sell => matching_core.bid,
1213 };
1214
1215 if released_price.is_none() {
1216 log::warn!(
1217 "Cannot release order {} yet: no market data available for {trigger_instrument_id}, will retry on next update",
1218 order.client_order_id(),
1219 );
1220 return None;
1221 }
1222
1223 released_price
1224 }
1225
1226 pub fn trigger_stop_order(&mut self, client_order_id: ClientOrderId) {
1230 let order = match self
1231 .cache
1232 .borrow()
1233 .order(&client_order_id)
1234 .map(|o| o.clone())
1235 {
1236 Some(order) => order,
1237 None => {
1238 log::error!(
1239 "Cannot trigger stop order: order {client_order_id} not found in cache"
1240 );
1241 return;
1242 }
1243 };
1244
1245 match order.order_type() {
1246 OrderType::StopLimit | OrderType::LimitIfTouched | OrderType::TrailingStopLimit => {
1247 self.fill_limit_order(client_order_id);
1248 }
1249 OrderType::Market
1250 | OrderType::MarketIfTouched
1251 | OrderType::StopMarket
1252 | OrderType::TrailingStopMarket => self.fill_market_order(client_order_id),
1253 _ => panic!("invalid `OrderType`, was {}", order.order_type()),
1254 }
1255 }
1256
1257 pub fn fill_limit_order(&mut self, client_order_id: ClientOrderId) {
1261 let order = match self
1262 .cache
1263 .borrow()
1264 .order(&client_order_id)
1265 .map(|o| o.clone())
1266 {
1267 Some(order) => order,
1268 None => {
1269 log::error!("Cannot fill limit order: order {client_order_id} not found in cache");
1270 return;
1271 }
1272 };
1273
1274 if matches!(order.order_type(), OrderType::Limit) {
1275 self.fill_market_order(client_order_id);
1276 return;
1277 }
1278
1279 let trigger_instrument_id = order
1280 .trigger_instrument_id()
1281 .unwrap_or(order.instrument_id());
1282
1283 let matching_core = match self.matching_cores.get(&trigger_instrument_id) {
1284 Some(core) => core,
1285 None => {
1286 log::error!(
1287 "Cannot fill limit order: no matching core for instrument {trigger_instrument_id}"
1288 );
1289 return; }
1291 };
1292
1293 let released_price =
1294 match self.validate_release(&order, matching_core, trigger_instrument_id) {
1295 Some(price) => price,
1296 None => return, };
1298
1299 let command = match self
1300 .manager
1301 .pop_submit_order_command(order.client_order_id())
1302 {
1303 Some(command) => command,
1304 None => return, };
1306
1307 if let Some(matching_core) = self.matching_cores.get_mut(&trigger_instrument_id) {
1308 if let Err(e) = matching_core.delete_order(client_order_id) {
1309 log::debug!("Error deleting order: {e:?}");
1310 }
1311
1312 let mut transformed = if let Ok(transformed) = LimitOrder::new_checked(
1314 order.trader_id(),
1315 order.strategy_id(),
1316 order.instrument_id(),
1317 order.client_order_id(),
1318 order.order_side(),
1319 order.quantity(),
1320 order.price().unwrap(),
1321 order.time_in_force(),
1322 order.expire_time(),
1323 order.is_post_only(),
1324 order.is_reduce_only(),
1325 order.is_quote_quantity(),
1326 order.display_qty(),
1327 None,
1328 Some(trigger_instrument_id),
1329 order.contingency_type(),
1330 order.order_list_id(),
1331 order.linked_order_ids().map(Vec::from),
1332 order.parent_order_id(),
1333 order.exec_algorithm_id(),
1334 order.exec_algorithm_params().cloned(),
1335 order.exec_spawn_id(),
1336 order.tags().map(Vec::from),
1337 UUID4::new(),
1338 self.clock.borrow().timestamp_ns(),
1339 ) {
1340 transformed
1341 } else {
1342 log::error!("Cannot create limit order");
1343 return;
1344 };
1345 transformed.liquidity_side = order.liquidity_side();
1346
1347 let original_events = order.events();
1354
1355 transformed.prepend_events(original_events.into_iter().cloned());
1356
1357 let add_result = {
1358 let mut cache = self.cache.borrow_mut();
1359 cache.add_order(
1360 OrderAny::Limit(transformed.clone()),
1361 command.position_id,
1362 command.client_id,
1363 true,
1364 )
1365 };
1366
1367 if let Err(e) = add_result {
1368 log::error!("Failed to add order: {e}");
1369 } else {
1370 msgbus::publish_order_event(
1371 format!("events.order.{}", order.strategy_id()).into(),
1372 transformed.last_event(),
1373 );
1374 }
1375
1376 let event = OrderReleased::new(
1377 order.trader_id(),
1378 order.strategy_id(),
1379 order.instrument_id(),
1380 order.client_order_id(),
1381 released_price,
1382 UUID4::new(),
1383 self.clock.borrow().timestamp_ns(),
1384 self.clock.borrow().timestamp_ns(),
1385 );
1386
1387 let event = OrderEventAny::Released(event);
1388
1389 let transformed = match self.cache.borrow_mut().update_order(&event) {
1390 Ok(order) => order,
1391 Err(e) => {
1392 log::error!("Failed to apply order event: {e}");
1393 return;
1394 }
1395 };
1396
1397 self.send_risk_event(event.clone());
1398
1399 log::info!("Releasing order {}", order.client_order_id());
1400
1401 msgbus::publish_order_event(
1403 format!("events.order.{}", transformed.strategy_id()).into(),
1404 &event,
1405 );
1406
1407 if let Some(exec_algorithm_id) = order.exec_algorithm_id() {
1408 self.send_algo_command(command, exec_algorithm_id);
1409 } else {
1410 self.send_exec_command(TradingCommand::SubmitOrder(command));
1411 }
1412 }
1413 }
1414
1415 pub fn fill_market_order(&mut self, client_order_id: ClientOrderId) {
1419 let mut order = match self
1420 .cache
1421 .borrow()
1422 .order(&client_order_id)
1423 .map(|o| o.clone())
1424 {
1425 Some(order) => order,
1426 None => {
1427 log::error!("Cannot fill market order: order {client_order_id} not found in cache");
1428 return;
1429 }
1430 };
1431
1432 let trigger_instrument_id = order
1433 .trigger_instrument_id()
1434 .unwrap_or(order.instrument_id());
1435
1436 let matching_core = match self.matching_cores.get(&trigger_instrument_id) {
1437 Some(core) => core,
1438 None => {
1439 log::error!(
1440 "Cannot fill market order: no matching core for instrument {trigger_instrument_id}"
1441 );
1442 return; }
1444 };
1445
1446 let released_price =
1447 match self.validate_release(&order, matching_core, trigger_instrument_id) {
1448 Some(price) => price,
1449 None => return, };
1451
1452 let command = self
1453 .manager
1454 .pop_submit_order_command(order.client_order_id())
1455 .expect("invalid operation `fill_market_order` with no command");
1456
1457 if let Some(matching_core) = self.matching_cores.get_mut(&trigger_instrument_id) {
1458 if let Err(e) = matching_core.delete_order(client_order_id) {
1459 log::debug!("Cannot delete order: {e:?}");
1460 }
1461
1462 order.set_emulation_trigger(None);
1463
1464 let mut transformed = MarketOrder::new(
1466 order.trader_id(),
1467 order.strategy_id(),
1468 order.instrument_id(),
1469 order.client_order_id(),
1470 order.order_side(),
1471 order.quantity(),
1472 order.time_in_force(),
1473 UUID4::new(),
1474 self.clock.borrow().timestamp_ns(),
1475 order.is_reduce_only(),
1476 order.is_quote_quantity(),
1477 order.contingency_type(),
1478 order.order_list_id(),
1479 order.linked_order_ids().map(Vec::from),
1480 order.parent_order_id(),
1481 order.exec_algorithm_id(),
1482 order.exec_algorithm_params().cloned(),
1483 order.exec_spawn_id(),
1484 order.tags().map(Vec::from),
1485 );
1486
1487 let original_events = order.events();
1488
1489 transformed.prepend_events(original_events.into_iter().cloned());
1490
1491 let add_result = {
1492 let mut cache = self.cache.borrow_mut();
1493 cache.add_order(
1494 OrderAny::Market(transformed.clone()),
1495 command.position_id,
1496 command.client_id,
1497 true,
1498 )
1499 };
1500
1501 if let Err(e) = add_result {
1502 log::error!("Failed to add order: {e}");
1503 } else {
1504 msgbus::publish_order_event(
1505 format!("events.order.{}", order.strategy_id()).into(),
1506 transformed.last_event(),
1507 );
1508 }
1509
1510 let ts_now = self.clock.borrow().timestamp_ns();
1511 let event = OrderReleased::new(
1512 order.trader_id(),
1513 order.strategy_id(),
1514 order.instrument_id(),
1515 order.client_order_id(),
1516 released_price,
1517 UUID4::new(),
1518 ts_now,
1519 ts_now,
1520 );
1521
1522 let event = OrderEventAny::Released(event);
1523
1524 if let Err(e) = self.cache.borrow_mut().update_order(&event) {
1525 log::error!("Failed to apply order event: {e}");
1526 return;
1527 }
1528 self.send_risk_event(event.clone());
1529
1530 log::info!("Releasing order {}", order.client_order_id());
1531
1532 msgbus::publish_order_event(
1534 format!("events.order.{}", order.strategy_id()).into(),
1535 &event,
1536 );
1537
1538 if let Some(exec_algorithm_id) = order.exec_algorithm_id() {
1539 self.send_algo_command(command, exec_algorithm_id);
1540 } else {
1541 self.send_exec_command(TradingCommand::SubmitOrder(command));
1542 }
1543 }
1544 }
1545
1546 fn update_trailing_stop_order(&mut self, order: &mut OrderAny) {
1547 let trigger_instrument_id = order
1548 .trigger_instrument_id()
1549 .unwrap_or_else(|| order.instrument_id());
1550 let Some(matching_core) = self.matching_cores.get(&trigger_instrument_id) else {
1551 log::error!(
1552 "Cannot update trailing-stop order: no matching core for instrument {trigger_instrument_id}"
1553 );
1554 return;
1555 };
1556
1557 let mut bid = matching_core.bid;
1558 let mut ask = matching_core.ask;
1559 let mut last = matching_core.last;
1560 let price_increment = matching_core.price_increment;
1561 let instrument_id = matching_core.instrument_id;
1562
1563 if bid.is_none() || ask.is_none() || last.is_none() {
1564 if let Some(q) = self.cache.borrow().quote(&instrument_id) {
1565 bid.get_or_insert(q.bid_price);
1566 ask.get_or_insert(q.ask_price);
1567 }
1568
1569 if let Some(t) = self.cache.borrow().trade(&instrument_id) {
1570 last.get_or_insert(t.price);
1571 }
1572 }
1573
1574 let was_activated = is_order_activated(order);
1575
1576 if !self.maybe_activate_trailing_stop(order, bid, ask, last) {
1577 return; }
1579
1580 if !was_activated {
1581 self.refresh_matching_core_entry(order, trigger_instrument_id);
1584 }
1585
1586 let (new_trigger_px, new_limit_px) = match trailing_stop_calculate(
1587 price_increment,
1588 order.trigger_price(),
1589 order,
1590 bid,
1591 ask,
1592 last,
1593 ) {
1594 Ok(pair) => pair,
1595 Err(e) => {
1596 log::warn!("Cannot calculate trailing-stop update: {e}");
1597 return;
1598 }
1599 };
1600
1601 if new_trigger_px.is_none() && new_limit_px.is_none() {
1602 return;
1603 }
1604
1605 let ts_now = self.clock.borrow().timestamp_ns();
1606 let update = OrderUpdated::new(
1607 order.trader_id(),
1608 order.strategy_id(),
1609 order.instrument_id(),
1610 order.client_order_id(),
1611 order.quantity(),
1612 UUID4::new(),
1613 ts_now,
1614 ts_now,
1615 false,
1616 order.venue_order_id(),
1617 order.account_id(),
1618 new_limit_px,
1619 new_trigger_px,
1620 None,
1621 order.is_quote_quantity(),
1622 );
1623 let wrapped = OrderEventAny::Updated(update);
1624
1625 *order = match self.cache.borrow_mut().update_order(&wrapped) {
1626 Ok(order) => order,
1627 Err(e) => {
1628 log::error!("Failed to apply order event: {e}");
1629 return;
1630 }
1631 };
1632
1633 self.refresh_matching_core_entry(order, trigger_instrument_id);
1634
1635 self.send_risk_event(wrapped);
1636 }
1637
1638 fn refresh_matching_core_entry(
1639 &mut self,
1640 order: &OrderAny,
1641 trigger_instrument_id: InstrumentId,
1642 ) {
1643 if let Some(matching_core) = self.matching_cores.get_mut(&trigger_instrument_id) {
1644 if let Err(e) = matching_core.delete_order(order.client_order_id()) {
1645 log::debug!("Cannot update trailing-stop match info: {e:?}");
1646 }
1647
1648 let (trigger_price, limit_price) = if order.trigger_price().is_some() {
1651 (order.trigger_price(), order.price())
1652 } else {
1653 (None, None)
1654 };
1655
1656 matching_core.add_order(RestingOrder::new(
1657 order.client_order_id(),
1658 order.order_side(),
1659 order.order_type(),
1660 trigger_price,
1661 limit_price,
1662 is_order_activated(order),
1663 ));
1664 }
1665 }
1666
1667 fn maybe_activate_trailing_stop(
1673 &self,
1674 order: &mut OrderAny,
1675 bid: Option<Price>,
1676 ask: Option<Price>,
1677 last: Option<Price>,
1678 ) -> bool {
1679 let (is_activated, activation_price, trigger_type, order_side) = match order {
1680 OrderAny::TrailingStopMarket(inner) => (
1681 inner.is_activated,
1682 inner.activation_price,
1683 inner.trigger_type,
1684 inner.order_side(),
1685 ),
1686 OrderAny::TrailingStopLimit(inner) => (
1687 inner.is_activated,
1688 inner.activation_price,
1689 inner.trigger_type,
1690 inner.order_side(),
1691 ),
1692 _ => return true,
1693 };
1694
1695 if is_activated {
1696 return true;
1697 }
1698
1699 if let Some(activation_price) = activation_price {
1700 let hit = match order_side {
1701 OrderSide::Buy => ask.is_some_and(|a| a <= activation_price),
1702 OrderSide::Sell => bid.is_some_and(|b| b >= activation_price),
1703 };
1704
1705 if hit {
1706 Self::set_trailing_stop_activated(order, None);
1707 self.persist_trailing_stop_activation(order);
1708 }
1709 return hit;
1710 }
1711
1712 let market_price = match trigger_type {
1713 TriggerType::LastPrice => last,
1714 _ => match order_side {
1715 OrderSide::Buy => ask,
1716 OrderSide::Sell => bid,
1717 },
1718 };
1719
1720 let Some(market_price) = market_price else {
1721 log::error!(
1722 "Cannot activate trailing stop {}: no market price available",
1723 order.client_order_id()
1724 );
1725 return false;
1726 };
1727
1728 Self::set_trailing_stop_activated(order, Some(market_price));
1729 self.persist_trailing_stop_activation(order);
1730 true
1731 }
1732
1733 fn set_trailing_stop_activated(order: &mut OrderAny, activation_price: Option<Price>) {
1734 match order {
1735 OrderAny::TrailingStopMarket(inner) => {
1736 if let Some(price) = activation_price {
1737 inner.activation_price = Some(price);
1738 }
1739 inner.set_activated();
1740 }
1741 OrderAny::TrailingStopLimit(inner) => {
1742 if let Some(price) = activation_price {
1743 inner.activation_price = Some(price);
1744 }
1745 inner.set_activated();
1746 }
1747 _ => {}
1748 }
1749 }
1750
1751 fn persist_trailing_stop_activation(&self, order: &OrderAny) {
1752 if let Err(e) = self.cache.borrow_mut().replace_order(order) {
1753 log::error!("Failed to update order: {e}");
1754 }
1755 }
1756
1757 fn send_algo_command(&self, command: SubmitOrder, exec_algorithm_id: ExecAlgorithmId) {
1758 let id = command.strategy_id;
1759 log::info!("{id} {CMD}{SEND} {command}");
1760
1761 let endpoint = format!("{exec_algorithm_id}.execute");
1762 msgbus::send_any(endpoint.into(), &TradingCommand::SubmitOrder(command));
1763 }
1764
1765 fn send_risk_command(&self, command: TradingCommand) {
1766 log_cmd_send(&command);
1767 let endpoint = MessagingSwitchboard::risk_engine_queue_execute();
1768 msgbus::send_trading_command(endpoint, command);
1769 }
1770
1771 fn send_exec_command(&self, command: TradingCommand) {
1772 log_cmd_send(&command);
1773 let endpoint = MessagingSwitchboard::exec_engine_queue_execute();
1774 msgbus::send_trading_command(endpoint, command);
1775 }
1776
1777 fn send_risk_event(&self, event: OrderEventAny) {
1778 log_evt_send(&event);
1779 let endpoint = MessagingSwitchboard::risk_engine_process();
1780 msgbus::send_order_event(endpoint, event);
1781 }
1782
1783 fn send_exec_event(&self, event: OrderEventAny) {
1784 log_evt_send(&event);
1785 let endpoint = MessagingSwitchboard::exec_engine_process();
1786 msgbus::send_order_event(endpoint, event);
1787 }
1788
1789 fn send_portfolio_order_event(&self, event: OrderEventAny) {
1790 log_evt_send(&event);
1791 let endpoint = MessagingSwitchboard::portfolio_update_order();
1792 msgbus::send_order_event(endpoint, event);
1793 }
1794
1795 fn send_data_command(&self, command: DataCommand) {
1796 log::info!("{CMD}{SEND} {command:?}");
1797 let endpoint = MessagingSwitchboard::data_engine_queue_execute();
1798 msgbus::send_data_command(endpoint, command);
1799 }
1800}
1801
1802fn publish_order_event(event: &OrderEventAny) {
1803 msgbus::publish_order_event(get_event_order_topic(event.strategy_id()), event);
1804
1805 if let OrderEventAny::Canceled(_) = event {
1806 msgbus::publish_order_event(get_order_canceled_topic(event.instrument_id()), event);
1807 }
1808}
1809
1810fn is_order_activated(order: &OrderAny) -> bool {
1811 match order {
1812 OrderAny::TrailingStopMarket(o) => o.is_activated,
1813 OrderAny::TrailingStopLimit(o) => o.is_activated,
1814 _ => true,
1815 }
1816}
1817
1818fn log_cmd_send(command: &TradingCommand) {
1819 if let Some(id) = command.strategy_id() {
1820 log::info!("{id} {CMD}{SEND} {command}");
1821 } else {
1822 log::info!("{CMD}{SEND} {command}");
1823 }
1824}
1825
1826fn log_evt_send(event: &OrderEventAny) {
1827 let id = event.strategy_id();
1828 log::info!("{id} {EVT}{SEND} {event}");
1829}
1830
1831#[cfg(test)]
1832mod tests {
1833 use std::{cell::RefCell, rc::Rc};
1834
1835 use nautilus_common::{
1836 cache::Cache,
1837 clock::VirtualClock,
1838 messages::data::{DataCommand, SubscribeCommand, UnsubscribeCommand},
1839 msgbus::{
1840 MessagingSwitchboard,
1841 stubs::{
1842 TypedIntoMessageSavingHandler, get_any_saving_handler,
1843 get_typed_into_message_saving_handler,
1844 },
1845 },
1846 };
1847 use nautilus_core::UUID4;
1848 use nautilus_model::{
1849 data::{QuoteTick, TradeTick},
1850 enums::{AggressorSide, OrderSide, OrderType, TrailingOffsetType, TriggerType},
1851 identifiers::{
1852 ClientId, ClientOrderId, OrderListId, StrategyId, Symbol, TradeId, TraderId,
1853 },
1854 instruments::{
1855 CryptoPerpetual, Instrument, InstrumentAny, SyntheticInstrument,
1856 stubs::crypto_perpetual_ethusdt,
1857 },
1858 orders::{OrderList, OrderTestBuilder},
1859 types::{Price, Quantity},
1860 };
1861 use rstest::{fixture, rstest};
1862 use rust_decimal_macros::dec;
1863 use ustr::Ustr;
1864
1865 use super::*;
1866
1867 #[fixture]
1868 fn instrument() -> CryptoPerpetual {
1869 crypto_perpetual_ethusdt()
1870 }
1871
1872 #[expect(clippy::type_complexity)]
1873 fn create_emulator() -> (
1874 Rc<RefCell<dyn Clock>>,
1875 Rc<RefCell<Cache>>,
1876 Rc<RefCell<OrderEmulator>>,
1877 ) {
1878 let clock: Rc<RefCell<dyn Clock>> = Rc::new(RefCell::new(VirtualClock::new()));
1879 let cache = Rc::new(RefCell::new(Cache::new(None, None)));
1880 let emulator = Rc::new(RefCell::new(OrderEmulator::new(
1881 clock.clone(),
1882 cache.clone(),
1883 )));
1884
1885 OrderEmulator::register_msgbus_handlers(&emulator);
1886
1887 (clock, cache, emulator)
1888 }
1889
1890 fn create_stop_market_order(instrument: &CryptoPerpetual, trigger: TriggerType) -> OrderAny {
1891 OrderTestBuilder::new(OrderType::StopMarket)
1892 .instrument_id(instrument.id())
1893 .side(OrderSide::Buy)
1894 .trigger_price(Price::from("5100.00"))
1895 .quantity(Quantity::from(1))
1896 .emulation_trigger(trigger)
1897 .build()
1898 }
1899
1900 fn create_stop_limit_order(instrument: &CryptoPerpetual, trigger: TriggerType) -> OrderAny {
1901 OrderTestBuilder::new(OrderType::StopLimit)
1902 .instrument_id(instrument.id())
1903 .side(OrderSide::Buy)
1904 .price(Price::from("5100.00"))
1905 .trigger_price(Price::from("5100.00"))
1906 .quantity(Quantity::from(1))
1907 .emulation_trigger(trigger)
1908 .build()
1909 }
1910
1911 fn create_list_stop_market_order(
1912 instrument: &CryptoPerpetual,
1913 client_order_id: &str,
1914 order_list_id: OrderListId,
1915 ) -> OrderAny {
1916 OrderTestBuilder::new(OrderType::StopMarket)
1917 .instrument_id(instrument.id())
1918 .client_order_id(ClientOrderId::from(client_order_id))
1919 .order_list_id(order_list_id)
1920 .side(OrderSide::Buy)
1921 .trigger_price(Price::from("5100.00"))
1922 .quantity(Quantity::from(1))
1923 .emulation_trigger(TriggerType::BidAsk)
1924 .build()
1925 }
1926
1927 fn create_submit_order(instrument: &CryptoPerpetual, order: &OrderAny) -> SubmitOrder {
1928 SubmitOrder::new(
1929 TraderId::from("TRADER-001"),
1930 None,
1931 StrategyId::from("STRATEGY-001"),
1932 instrument.id(),
1933 order.client_order_id(),
1934 order.init_event().clone(),
1935 None,
1936 None,
1937 None,
1938 UUID4::new(),
1939 0.into(),
1940 None, )
1942 }
1943
1944 fn create_quote_tick(instrument: &CryptoPerpetual, bid: &str, ask: &str) -> QuoteTick {
1945 QuoteTick::new(
1946 instrument.id(),
1947 Price::from(bid),
1948 Price::from(ask),
1949 Quantity::from(10),
1950 Quantity::from(10),
1951 0.into(),
1952 0.into(),
1953 )
1954 }
1955
1956 fn create_trade_tick(instrument: &CryptoPerpetual, price: &str) -> TradeTick {
1957 TradeTick::new(
1958 instrument.id(),
1959 Price::from(price),
1960 Quantity::from(1),
1961 AggressorSide::Buy,
1962 TradeId::from("T-001"),
1963 0.into(),
1964 0.into(),
1965 )
1966 }
1967
1968 fn add_instrument_to_cache(cache: &Rc<RefCell<Cache>>, instrument: &CryptoPerpetual) {
1969 cache
1970 .borrow_mut()
1971 .add_instrument(InstrumentAny::CryptoPerpetual(instrument.clone()))
1972 .unwrap();
1973 }
1974
1975 fn register_risk_event_handler(id: &str) -> TypedIntoMessageSavingHandler<OrderEventAny> {
1976 let (handler, saving_handler) =
1977 get_typed_into_message_saving_handler::<OrderEventAny>(Some(Ustr::from(id)));
1978 msgbus::register_order_event_endpoint(MessagingSwitchboard::risk_engine_process(), handler);
1979 saving_handler
1980 }
1981
1982 fn register_exec_event_handler(
1983 cache: Rc<RefCell<Cache>>,
1984 id: &str,
1985 ) -> TypedIntoMessageSavingHandler<OrderEventAny> {
1986 let messages = Rc::new(RefCell::new(Vec::new()));
1987 let messages_for_handler = messages.clone();
1988 msgbus::register_order_event_endpoint(
1989 MessagingSwitchboard::exec_engine_process(),
1990 TypedIntoHandler::from(move |event: OrderEventAny| {
1991 cache.borrow_mut().update_order(&event).unwrap();
1992 messages_for_handler.borrow_mut().push(event);
1993 }),
1994 );
1995 TypedIntoMessageSavingHandler::new_with_messages(Some(Ustr::from(id)), messages)
1996 }
1997
1998 fn register_portfolio_event_handler(id: &str) -> TypedIntoMessageSavingHandler<OrderEventAny> {
1999 let (handler, saving_handler) =
2000 get_typed_into_message_saving_handler::<OrderEventAny>(Some(Ustr::from(id)));
2001 msgbus::register_order_event_endpoint(
2002 MessagingSwitchboard::portfolio_update_order(),
2003 handler,
2004 );
2005 saving_handler
2006 }
2007
2008 fn register_data_command_handler(id: &str) -> TypedIntoMessageSavingHandler<DataCommand> {
2009 let (handler, saving_handler) =
2010 get_typed_into_message_saving_handler::<DataCommand>(Some(Ustr::from(id)));
2011 msgbus::register_data_command_endpoint(
2012 MessagingSwitchboard::data_engine_queue_execute(),
2013 handler,
2014 );
2015 saving_handler
2016 }
2017
2018 fn subscribe_order_topic(
2019 strategy_id: StrategyId,
2020 ) -> (TypedHandler<OrderEventAny>, Rc<RefCell<Vec<OrderEventAny>>>) {
2021 let events = Rc::new(RefCell::new(Vec::new()));
2022 let handler = TypedHandler::from({
2023 let events = events.clone();
2024 move |event: &OrderEventAny| {
2025 events.borrow_mut().push(event.clone());
2026 }
2027 });
2028 msgbus::subscribe_order_events(
2029 format!("events.order.{strategy_id}").into(),
2030 handler.clone(),
2031 None,
2032 );
2033 (handler, events)
2034 }
2035
2036 fn subscribe_order_cancel_topic(
2037 instrument_id: InstrumentId,
2038 ) -> (TypedHandler<OrderEventAny>, Rc<RefCell<Vec<OrderEventAny>>>) {
2039 let events = Rc::new(RefCell::new(Vec::new()));
2040 let handler = TypedHandler::from({
2041 let events = events.clone();
2042 move |event: &OrderEventAny| {
2043 events.borrow_mut().push(event.clone());
2044 }
2045 });
2046 msgbus::subscribe_order_events(
2047 get_order_canceled_topic(instrument_id).into(),
2048 handler.clone(),
2049 None,
2050 );
2051 (handler, events)
2052 }
2053
2054 #[rstest]
2055 fn test_dispatch_manager_publish_initialized_publishes_order_event(
2056 instrument: CryptoPerpetual,
2057 ) {
2058 let (_clock, _cache, emulator) = create_emulator();
2059 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2060 let strategy_id = order.strategy_id();
2061 let client_order_id = order.client_order_id();
2062 let event = OrderEventAny::Initialized(order.init_event().clone());
2063 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
2064
2065 emulator
2066 .borrow_mut()
2067 .dispatch_manager_action(OrderManagerAction::PublishInitialized(event));
2068 msgbus::unsubscribe_order_events(
2069 format!("events.order.{strategy_id}").into(),
2070 &order_handler,
2071 );
2072 let order_events = order_events.borrow();
2073
2074 assert_eq!(order_events.len(), 1);
2075 assert!(matches!(
2076 &order_events[0],
2077 OrderEventAny::Initialized(event) if event.client_order_id == client_order_id
2078 ));
2079 }
2080
2081 #[rstest]
2082 fn test_dispatch_manager_submit_to_emulator_applies_emulated_event(
2083 instrument: CryptoPerpetual,
2084 ) {
2085 let (_clock, cache, emulator) = create_emulator();
2086 let risk_events = register_risk_event_handler("RiskEngine.process.dispatch_emulated");
2087 add_instrument_to_cache(&cache, &instrument);
2088 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2089 let client_order_id = order.client_order_id();
2090 let command = create_submit_order(&instrument, &order);
2091 cache
2092 .borrow_mut()
2093 .add_order(order, None, None, false)
2094 .unwrap();
2095 emulator
2096 .borrow_mut()
2097 .cache_submit_order_command(command.clone());
2098
2099 emulator
2100 .borrow_mut()
2101 .dispatch_manager_action(OrderManagerAction::SubmitToEmulator(command));
2102 let cache = cache.borrow();
2103 let cached_order = cache.order(&client_order_id).unwrap();
2104 let risk_events = risk_events.get_messages();
2105
2106 assert_eq!(cached_order.status(), OrderStatus::Emulated);
2107 assert_eq!(risk_events.len(), 1);
2108 assert!(matches!(risk_events[0], OrderEventAny::Emulated(_)));
2109 }
2110
2111 #[rstest]
2112 fn test_registered_execute_endpoint_routes_submit_order(instrument: CryptoPerpetual) {
2113 let (_clock, cache, emulator) = create_emulator();
2114 let risk_events = register_risk_event_handler("RiskEngine.process.endpoint_emulated");
2115 let data_commands = register_data_command_handler("DataEngine.queue_execute.endpoint");
2116 add_instrument_to_cache(&cache, &instrument);
2117 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2118 let client_order_id = order.client_order_id();
2119 let command = create_submit_order(&instrument, &order);
2120 cache
2121 .borrow_mut()
2122 .add_order(order, None, None, false)
2123 .unwrap();
2124 emulator
2125 .borrow_mut()
2126 .cache_submit_order_command(command.clone());
2127
2128 msgbus::send_trading_command(
2129 MessagingSwitchboard::order_emulator_execute(),
2130 TradingCommand::SubmitOrder(command),
2131 );
2132 let cache = cache.borrow();
2133 let cached_order = cache.order(&client_order_id).unwrap();
2134 let risk_events = risk_events.get_messages();
2135 let data_commands = data_commands.get_messages();
2136
2137 assert_eq!(cached_order.status(), OrderStatus::Emulated);
2138 assert_eq!(risk_events.len(), 1);
2139 assert!(matches!(risk_events[0], OrderEventAny::Emulated(_)));
2140 assert_eq!(data_commands.len(), 1);
2141 assert!(matches!(
2142 &data_commands[0],
2143 DataCommand::Subscribe(SubscribeCommand::Quotes(command))
2144 if command.instrument_id == instrument.id()
2145 ));
2146 }
2147
2148 #[rstest]
2149 fn test_cancel_all_orders_filters_emulated_orders_by_client(instrument: CryptoPerpetual) {
2150 let (_clock, cache, emulator) = create_emulator();
2151 let _risk_events = register_risk_event_handler("RiskEngine.process.cancel_all_client");
2152 let portfolio_events =
2153 register_portfolio_event_handler("Portfolio.update_order.cancel_all_client");
2154 let _data_commands = register_data_command_handler("DataEngine.queue_execute.cancel_all");
2155 add_instrument_to_cache(&cache, &instrument);
2156
2157 let selected_client = ClientId::from("CLIENT-001");
2158 let other_client = ClientId::from("CLIENT-002");
2159 let make_order = |client_order_id| {
2160 OrderTestBuilder::new(OrderType::StopMarket)
2161 .instrument_id(instrument.id())
2162 .client_order_id(client_order_id)
2163 .side(OrderSide::Buy)
2164 .trigger_price(Price::from("5100.00"))
2165 .quantity(Quantity::from(1))
2166 .emulation_trigger(TriggerType::BidAsk)
2167 .build()
2168 };
2169 let selected_order = make_order(ClientOrderId::from("O-EMULATED-SELECTED"));
2170 let other_order = make_order(ClientOrderId::from("O-EMULATED-OTHER"));
2171 let unclaimed_order = make_order(ClientOrderId::from("O-EMULATED-UNCLAIMED"));
2172 let synthetic_formula = format!("{} * 1.0", instrument.id());
2173 let synthetic = SyntheticInstrument::builder()
2174 .symbol(Symbol::from("ETH-INDEX"))
2175 .price_precision(instrument.price_precision())
2176 .components(vec![instrument.id()])
2177 .formula(&synthetic_formula)
2178 .ts_event(0.into())
2179 .ts_init(0.into())
2180 .build()
2181 .unwrap();
2182 let trigger_instrument_id = synthetic.id;
2183 let cross_trigger_order = OrderTestBuilder::new(OrderType::StopMarket)
2184 .instrument_id(instrument.id())
2185 .client_order_id(ClientOrderId::from("O-EMULATED-CROSS-TRIGGER"))
2186 .side(OrderSide::Buy)
2187 .trigger_price(Price::from("5100.00"))
2188 .quantity(Quantity::from(1))
2189 .emulation_trigger(TriggerType::BidAsk)
2190 .trigger_instrument_id(trigger_instrument_id)
2191 .build();
2192 let mut selected_submit = create_submit_order(&instrument, &selected_order);
2193 selected_submit.client_id = Some(selected_client);
2194 let mut other_submit = create_submit_order(&instrument, &other_order);
2195 other_submit.client_id = Some(other_client);
2196 let unclaimed_submit = create_submit_order(&instrument, &unclaimed_order);
2197 let mut cross_trigger_submit = create_submit_order(&instrument, &cross_trigger_order);
2198 cross_trigger_submit.client_id = Some(selected_client);
2199 {
2200 let mut cache = cache.borrow_mut();
2201 cache.add_synthetic(synthetic).unwrap();
2202 cache
2203 .add_order(selected_order.clone(), None, Some(selected_client), false)
2204 .unwrap();
2205 cache
2206 .add_order(other_order.clone(), None, Some(other_client), false)
2207 .unwrap();
2208 cache
2209 .add_order(unclaimed_order.clone(), None, None, false)
2210 .unwrap();
2211 cache
2212 .add_order(
2213 cross_trigger_order.clone(),
2214 None,
2215 Some(selected_client),
2216 false,
2217 )
2218 .unwrap();
2219 }
2220 emulator.borrow_mut().handle_submit_order(&selected_submit);
2221 emulator.borrow_mut().handle_submit_order(&other_submit);
2222 emulator.borrow_mut().handle_submit_order(&unclaimed_submit);
2223 emulator
2224 .borrow_mut()
2225 .handle_submit_order(&cross_trigger_submit);
2226
2227 emulator
2228 .borrow_mut()
2229 .execute(TradingCommand::CancelAllOrders(CancelAllOrders::new(
2230 TraderId::from("TRADER-001"),
2231 Some(selected_client),
2232 StrategyId::from("CALLER-001"),
2233 trigger_instrument_id,
2234 None,
2235 UUID4::new(),
2236 0.into(),
2237 None,
2238 None,
2239 )));
2240
2241 assert_eq!(
2242 cache
2243 .borrow()
2244 .order(&cross_trigger_order.client_order_id())
2245 .unwrap()
2246 .status(),
2247 OrderStatus::Emulated
2248 );
2249
2250 emulator
2251 .borrow_mut()
2252 .execute(TradingCommand::CancelAllOrders(CancelAllOrders::new(
2253 TraderId::from("TRADER-001"),
2254 Some(selected_client),
2255 StrategyId::from("CALLER-001"),
2256 instrument.id(),
2257 None,
2258 UUID4::new(),
2259 0.into(),
2260 None,
2261 None,
2262 )));
2263
2264 let cache = cache.borrow();
2265 assert_eq!(
2266 cache
2267 .order(&selected_order.client_order_id())
2268 .unwrap()
2269 .status(),
2270 OrderStatus::Canceled
2271 );
2272 assert_eq!(
2273 cache
2274 .order(&other_order.client_order_id())
2275 .unwrap()
2276 .status(),
2277 OrderStatus::Emulated
2278 );
2279 assert_eq!(
2280 cache
2281 .order(&unclaimed_order.client_order_id())
2282 .unwrap()
2283 .status(),
2284 OrderStatus::Emulated
2285 );
2286 assert_eq!(
2287 cache
2288 .order(&cross_trigger_order.client_order_id())
2289 .unwrap()
2290 .status(),
2291 OrderStatus::Canceled
2292 );
2293 let portfolio_events = portfolio_events.get_messages();
2294 assert_eq!(portfolio_events.len(), 2);
2295 let canceled_ids: AHashSet<_> = portfolio_events
2296 .iter()
2297 .filter_map(|event| match event {
2298 OrderEventAny::Canceled(event) => Some(event.client_order_id),
2299 _ => None,
2300 })
2301 .collect();
2302 assert_eq!(
2303 canceled_ids,
2304 AHashSet::from_iter([
2305 selected_order.client_order_id(),
2306 cross_trigger_order.client_order_id(),
2307 ])
2308 );
2309 }
2310
2311 #[rstest]
2312 fn test_dispatch_manager_submit_to_risk_uses_risk_queue(instrument: CryptoPerpetual) {
2313 let (_clock, _cache, emulator) = create_emulator();
2314 let (handler, messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
2315 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
2316 msgbus::register_trading_command_endpoint(
2317 MessagingSwitchboard::risk_engine_queue_execute(),
2318 handler,
2319 );
2320 let order = create_stop_market_order(&instrument, TriggerType::Default);
2321 let command = create_submit_order(&instrument, &order);
2322 let client_order_id = command.client_order_id;
2323
2324 emulator
2325 .borrow_mut()
2326 .dispatch_manager_action(OrderManagerAction::SubmitToRisk(command));
2327
2328 let messages = messages.get_messages();
2329 assert_eq!(messages.len(), 1);
2330 assert!(matches!(
2331 messages.first(),
2332 Some(TradingCommand::SubmitOrder(command))
2333 if command.client_order_id == client_order_id
2334 ));
2335 }
2336
2337 #[rstest]
2338 fn test_dispatch_manager_submit_to_algorithm_uses_dynamic_endpoint(
2339 instrument: CryptoPerpetual,
2340 ) {
2341 let (_clock, _cache, emulator) = create_emulator();
2342 let exec_algorithm_id = ExecAlgorithmId::from("ALG-001");
2343 let endpoint = format!("{exec_algorithm_id}.execute");
2344 let (handler, messages) =
2345 get_any_saving_handler::<TradingCommand>(Some(Ustr::from("ALG-001.execute")));
2346 msgbus::register_any(endpoint.into(), handler);
2347 let order = create_stop_market_order(&instrument, TriggerType::Default);
2348 let command = create_submit_order(&instrument, &order);
2349 let client_order_id = command.client_order_id;
2350
2351 emulator
2352 .borrow_mut()
2353 .dispatch_manager_action(OrderManagerAction::SubmitToAlgorithm {
2354 command,
2355 exec_algorithm_id,
2356 });
2357
2358 let messages = messages.get_messages();
2359 assert_eq!(messages.len(), 1);
2360 assert!(matches!(
2361 messages.first(),
2362 Some(TradingCommand::SubmitOrder(command))
2363 if command.client_order_id == client_order_id
2364 ));
2365 }
2366
2367 #[rstest]
2368 fn test_dispatch_manager_cancel_local_applies_and_publishes_event(instrument: CryptoPerpetual) {
2369 let (_clock, cache, emulator) = create_emulator();
2370 let portfolio_events =
2371 register_portfolio_event_handler("Portfolio.update_order.dispatch_canceled");
2372 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2373 let client_order_id = order.client_order_id();
2374 let strategy_id = order.strategy_id();
2375 let instrument_id = order.instrument_id();
2376 let command = create_submit_order(&instrument, &order);
2377 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
2378 let (cancel_handler, cancel_events) = subscribe_order_cancel_topic(instrument_id);
2379 cache
2380 .borrow_mut()
2381 .add_order(order.clone(), None, None, false)
2382 .unwrap();
2383 emulator.borrow_mut().cache_submit_order_command(command);
2384
2385 emulator
2386 .borrow_mut()
2387 .dispatch_manager_action(OrderManagerAction::CancelLocal(order));
2388 msgbus::unsubscribe_order_events(get_event_order_topic(strategy_id).into(), &order_handler);
2389 msgbus::unsubscribe_order_events(
2390 get_order_canceled_topic(instrument_id).into(),
2391 &cancel_handler,
2392 );
2393 let cached_order = cache
2394 .borrow()
2395 .order(&client_order_id)
2396 .map(|order| order.clone())
2397 .unwrap();
2398 let portfolio_events = portfolio_events.get_messages();
2399 let order_events = order_events.borrow();
2400 let cancel_events = cancel_events.borrow();
2401 let commands = emulator.borrow().get_submit_order_commands();
2402
2403 assert_eq!(cached_order.status(), OrderStatus::Canceled);
2404 assert_eq!(portfolio_events.len(), 1);
2405 assert_eq!(order_events.len(), 1);
2406 assert_eq!(cancel_events.len(), 1);
2407 assert!(matches!(
2408 &portfolio_events[0],
2409 OrderEventAny::Canceled(event) if event.client_order_id == client_order_id
2410 ));
2411 assert_eq!(order_events[0].client_order_id(), client_order_id);
2412 assert_eq!(cancel_events[0].client_order_id(), client_order_id);
2413 assert!(!commands.contains_key(&client_order_id));
2414 }
2415
2416 #[rstest]
2417 fn test_dispatch_manager_modify_local_quantity_updates_order(instrument: CryptoPerpetual) {
2418 let (_clock, cache, emulator) = create_emulator();
2419 let risk_events = register_risk_event_handler("RiskEngine.process.dispatch_updated");
2420 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2421 let client_order_id = order.client_order_id();
2422 let quantity = Quantity::from(2);
2423 cache
2424 .borrow_mut()
2425 .add_order(order.clone(), None, None, false)
2426 .unwrap();
2427
2428 emulator
2429 .borrow_mut()
2430 .dispatch_manager_action(OrderManagerAction::ModifyLocalQuantity { order, quantity });
2431 let cache = cache.borrow();
2432 let cached_order = cache.order(&client_order_id).unwrap();
2433 let risk_events = risk_events.get_messages();
2434
2435 assert_eq!(cached_order.quantity(), quantity);
2436 assert_eq!(risk_events.len(), 1);
2437 assert!(matches!(
2438 &risk_events[0],
2439 OrderEventAny::Updated(event) if event.client_order_id == client_order_id
2440 ));
2441 }
2442
2443 #[rstest]
2444 fn test_subscribed_quotes_initially_empty() {
2445 let (_clock, _cache, emulator) = create_emulator();
2446
2447 assert!(emulator.borrow().subscribed_quotes().is_empty());
2448 }
2449
2450 #[rstest]
2451 fn test_subscribed_trades_initially_empty() {
2452 let (_clock, _cache, emulator) = create_emulator();
2453
2454 assert!(emulator.borrow().subscribed_trades().is_empty());
2455 }
2456
2457 #[rstest]
2458 fn test_get_submit_order_commands_initially_empty() {
2459 let (_clock, _cache, emulator) = create_emulator();
2460
2461 assert!(emulator.borrow().get_submit_order_commands().is_empty());
2462 }
2463
2464 #[rstest]
2465 fn test_get_matching_core_returns_none_when_not_created(instrument: CryptoPerpetual) {
2466 let (_clock, _cache, emulator) = create_emulator();
2467
2468 assert!(
2469 emulator
2470 .borrow()
2471 .get_matching_core(&instrument.id())
2472 .is_none()
2473 );
2474 }
2475
2476 #[rstest]
2477 fn test_create_matching_core(instrument: CryptoPerpetual) {
2478 let (_clock, _cache, emulator) = create_emulator();
2479
2480 emulator
2481 .borrow_mut()
2482 .create_matching_core(instrument.id(), instrument.price_increment);
2483
2484 assert!(
2485 emulator
2486 .borrow()
2487 .get_matching_core(&instrument.id())
2488 .is_some()
2489 );
2490 }
2491
2492 #[rstest]
2493 fn test_on_quote_tick_no_matching_core_does_not_panic(instrument: CryptoPerpetual) {
2494 let (_clock, _cache, emulator) = create_emulator();
2495 let quote = create_quote_tick(&instrument, "5060.00", "5070.00");
2496
2497 emulator.borrow_mut().on_quote_tick(quote);
2498 }
2499
2500 #[rstest]
2501 fn test_on_trade_tick_no_matching_core_does_not_panic(instrument: CryptoPerpetual) {
2502 let (_clock, _cache, emulator) = create_emulator();
2503 let trade = create_trade_tick(&instrument, "5065.00");
2504
2505 emulator.borrow_mut().on_trade_tick(trade);
2506 }
2507
2508 #[rstest]
2509 fn test_submit_order_bid_ask_trigger_creates_matching_core(instrument: CryptoPerpetual) {
2510 let (_clock, cache, emulator) = create_emulator();
2511 add_instrument_to_cache(&cache, &instrument);
2512 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2513 let command = create_submit_order(&instrument, &order);
2514 cache
2515 .borrow_mut()
2516 .add_order(order, None, None, false)
2517 .unwrap();
2518
2519 emulator
2520 .borrow_mut()
2521 .cache_submit_order_command(command.clone());
2522 emulator.borrow_mut().handle_submit_order(&command);
2523
2524 assert!(
2525 emulator
2526 .borrow()
2527 .get_matching_core(&instrument.id())
2528 .is_some()
2529 );
2530 }
2531
2532 #[rstest]
2533 fn test_submit_order_bid_ask_trigger_tracks_quote_subscription(instrument: CryptoPerpetual) {
2534 let (_clock, cache, emulator) = create_emulator();
2535 let data_commands = register_data_command_handler("DataEngine.queue_execute.quotes");
2536 add_instrument_to_cache(&cache, &instrument);
2537 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2538 let command = create_submit_order(&instrument, &order);
2539 cache
2540 .borrow_mut()
2541 .add_order(order, None, None, false)
2542 .unwrap();
2543
2544 emulator
2545 .borrow_mut()
2546 .cache_submit_order_command(command.clone());
2547 emulator.borrow_mut().handle_submit_order(&command);
2548
2549 assert_eq!(emulator.borrow().subscribed_quotes(), vec![instrument.id()]);
2550 assert!(emulator.borrow().subscribed_trades().is_empty());
2551
2552 let commands = data_commands.get_messages();
2553 assert_eq!(commands.len(), 1);
2554 assert!(matches!(
2555 &commands[0],
2556 DataCommand::Subscribe(SubscribeCommand::Quotes(command))
2557 if command.instrument_id == instrument.id()
2558 ));
2559
2560 let quote = create_quote_tick(&instrument, "5060.00", "5070.00");
2561 msgbus::publish_quote(get_quotes_topic(instrument.id()), "e);
2562 let core = emulator
2563 .borrow()
2564 .get_matching_core(&instrument.id())
2565 .unwrap();
2566 assert_eq!(core.bid, Some(Price::from("5060.00")));
2567 assert_eq!(core.ask, Some(Price::from("5070.00")));
2568 }
2569
2570 #[rstest]
2571 fn test_submit_order_last_price_trigger_tracks_trade_subscription(instrument: CryptoPerpetual) {
2572 let (_clock, cache, emulator) = create_emulator();
2573 let data_commands = register_data_command_handler("DataEngine.queue_execute.trades");
2574 add_instrument_to_cache(&cache, &instrument);
2575 let order = create_stop_market_order(&instrument, TriggerType::LastPrice);
2576 let command = create_submit_order(&instrument, &order);
2577 cache
2578 .borrow_mut()
2579 .add_order(order, None, None, false)
2580 .unwrap();
2581
2582 emulator
2583 .borrow_mut()
2584 .cache_submit_order_command(command.clone());
2585 emulator.borrow_mut().handle_submit_order(&command);
2586
2587 assert!(emulator.borrow().subscribed_quotes().is_empty());
2588 assert_eq!(emulator.borrow().subscribed_trades(), vec![instrument.id()]);
2589
2590 let commands = data_commands.get_messages();
2591 assert_eq!(commands.len(), 1);
2592 assert!(matches!(
2593 &commands[0],
2594 DataCommand::Subscribe(SubscribeCommand::Trades(command))
2595 if command.instrument_id == instrument.id()
2596 ));
2597
2598 let trade = create_trade_tick(&instrument, "5065.00");
2599 msgbus::publish_trade(get_trades_topic(instrument.id()), &trade);
2600 let core = emulator
2601 .borrow()
2602 .get_matching_core(&instrument.id())
2603 .unwrap();
2604 assert_eq!(core.last, Some(Price::from("5065.00")));
2605 }
2606
2607 #[rstest]
2608 fn test_reset_unsubscribes_market_data_and_clears_state(instrument: CryptoPerpetual) {
2609 let (_clock, cache, emulator) = create_emulator();
2610 let data_commands = register_data_command_handler("DataEngine.queue_execute.reset");
2611 add_instrument_to_cache(&cache, &instrument);
2612 let quote_order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2613 let quote_command = create_submit_order(&instrument, "e_order);
2614 let trade_order = OrderTestBuilder::new(OrderType::StopMarket)
2615 .instrument_id(instrument.id())
2616 .client_order_id(ClientOrderId::from("O-RESET-TRADE"))
2617 .side(OrderSide::Buy)
2618 .trigger_price(Price::from("5100.00"))
2619 .quantity(Quantity::from(1))
2620 .emulation_trigger(TriggerType::LastPrice)
2621 .build();
2622 let trade_command = create_submit_order(&instrument, &trade_order);
2623 cache
2624 .borrow_mut()
2625 .add_order(quote_order, None, None, false)
2626 .unwrap();
2627 cache
2628 .borrow_mut()
2629 .add_order(trade_order, None, None, false)
2630 .unwrap();
2631 emulator
2632 .borrow_mut()
2633 .cache_submit_order_command(quote_command.clone());
2634 emulator.borrow_mut().handle_submit_order("e_command);
2635 emulator
2636 .borrow_mut()
2637 .cache_submit_order_command(trade_command.clone());
2638 emulator.borrow_mut().handle_submit_order(&trade_command);
2639 data_commands.clear();
2640
2641 emulator.borrow_mut().reset();
2642 let commands = data_commands.get_messages();
2643 let emulator_ref = emulator.borrow();
2644
2645 assert!(emulator_ref.subscribed_quotes.is_empty());
2646 assert!(emulator_ref.subscribed_trades.is_empty());
2647 assert!(emulator_ref.subscribed_strategies.is_empty());
2648 assert!(emulator_ref.monitored_positions.is_empty());
2649 assert!(emulator_ref.quote_handlers.is_empty());
2650 assert!(emulator_ref.trade_handlers.is_empty());
2651 assert!(emulator_ref.get_submit_order_commands().is_empty());
2652 assert!(emulator_ref.get_matching_core(&instrument.id()).is_none());
2653 assert!(commands.iter().any(|command| matches!(
2654 command,
2655 DataCommand::Unsubscribe(UnsubscribeCommand::Quotes(command))
2656 if command.instrument_id == instrument.id()
2657 )));
2658 assert!(commands.iter().any(|command| matches!(
2659 command,
2660 DataCommand::Unsubscribe(UnsubscribeCommand::Trades(command))
2661 if command.instrument_id == instrument.id()
2662 )));
2663
2664 drop(emulator_ref);
2665 emulator
2666 .borrow_mut()
2667 .create_matching_core(instrument.id(), instrument.price_increment);
2668 let quote = create_quote_tick(&instrument, "5060.00", "5070.00");
2669 let trade = create_trade_tick(&instrument, "5065.00");
2670 msgbus::publish_quote(get_quotes_topic(instrument.id()), "e);
2671 msgbus::publish_trade(get_trades_topic(instrument.id()), &trade);
2672 let core = emulator
2673 .borrow()
2674 .get_matching_core(&instrument.id())
2675 .unwrap();
2676 assert_eq!(core.bid, None);
2677 assert_eq!(core.ask, None);
2678 assert_eq!(core.last, None);
2679 }
2680
2681 #[rstest]
2682 fn test_reset_unsubscribes_in_sorted_instrument_order(instrument: CryptoPerpetual) {
2683 let (_clock, cache, emulator) = create_emulator();
2684 let data_commands = register_data_command_handler("DataEngine.queue_execute.reset_order");
2685
2686 let ids = [
2689 "SOLUSDT-PERP.BINANCE",
2690 "ADAUSDT-PERP.BINANCE",
2691 "XRPUSDT-PERP.BINANCE",
2692 "BTCUSDT-PERP.BINANCE",
2693 "DOTUSDT-PERP.BINANCE",
2694 ];
2695
2696 for (i, id) in ids.iter().enumerate() {
2697 let mut variant = instrument.clone();
2698 variant.id = InstrumentId::from(*id);
2699 add_instrument_to_cache(&cache, &variant);
2700
2701 let order = OrderTestBuilder::new(OrderType::StopMarket)
2702 .instrument_id(variant.id)
2703 .client_order_id(ClientOrderId::from(format!("O-RESET-{i}").as_str()))
2704 .side(OrderSide::Buy)
2705 .trigger_price(Price::from("5100.00"))
2706 .quantity(Quantity::from(1))
2707 .emulation_trigger(TriggerType::BidAsk)
2708 .build();
2709 let command = create_submit_order(&variant, &order);
2710 cache
2711 .borrow_mut()
2712 .add_order(order, None, None, false)
2713 .unwrap();
2714 emulator
2715 .borrow_mut()
2716 .cache_submit_order_command(command.clone());
2717 emulator.borrow_mut().handle_submit_order(&command);
2718 }
2719 data_commands.clear();
2720
2721 emulator.borrow_mut().reset();
2722
2723 let unsubscribed: Vec<InstrumentId> = data_commands
2724 .get_messages()
2725 .iter()
2726 .filter_map(|command| match command {
2727 DataCommand::Unsubscribe(UnsubscribeCommand::Quotes(command)) => {
2728 Some(command.instrument_id)
2729 }
2730 _ => None,
2731 })
2732 .collect();
2733
2734 let mut expected: Vec<InstrumentId> =
2735 ids.iter().map(|id| InstrumentId::from(*id)).collect();
2736 expected.sort();
2737
2738 assert_eq!(unsubscribed, expected);
2739 }
2740
2741 #[rstest]
2742 fn test_submit_order_list_handles_reentrant_order_events(instrument: CryptoPerpetual) {
2743 let (_clock, cache, emulator) = create_emulator();
2744 let data_commands = register_data_command_handler("DataEngine.queue_execute.list");
2745 add_instrument_to_cache(&cache, &instrument);
2746 let order_list_id = OrderListId::from("OL-EMULATOR-001");
2747 let first_order = create_list_stop_market_order(&instrument, "O-LIST-001", order_list_id);
2748 let second_order = create_list_stop_market_order(&instrument, "O-LIST-002", order_list_id);
2749 let orders = vec![first_order.clone(), second_order.clone()];
2750 let order_list = OrderList::from_orders(&orders, 0.into());
2751 let order_inits = orders
2752 .iter()
2753 .map(|order| order.init_event().clone())
2754 .collect();
2755 cache
2756 .borrow_mut()
2757 .add_order(first_order.clone(), None, None, false)
2758 .unwrap();
2759 cache
2760 .borrow_mut()
2761 .add_order(second_order.clone(), None, None, false)
2762 .unwrap();
2763 let command = SubmitOrderList::new(
2764 TraderId::from("TRADER-001"),
2765 None,
2766 StrategyId::from("STRATEGY-001"),
2767 order_list,
2768 order_inits,
2769 None,
2770 None,
2771 None,
2772 UUID4::new(),
2773 0.into(),
2774 None,
2775 );
2776
2777 emulator
2778 .borrow_mut()
2779 .execute(TradingCommand::SubmitOrderList(command));
2780
2781 let commands = data_commands.get_messages();
2782 let cache = cache.borrow();
2783 let first_status = cache
2784 .order(&first_order.client_order_id())
2785 .unwrap()
2786 .status();
2787 let second_status = cache
2788 .order(&second_order.client_order_id())
2789 .unwrap()
2790 .status();
2791 drop(cache);
2792 let emulator = emulator.borrow();
2793
2794 assert_eq!(first_status, OrderStatus::Emulated);
2795 assert_eq!(second_status, OrderStatus::Emulated);
2796 assert!(emulator.get_matching_core(&instrument.id()).is_some());
2797 assert_eq!(emulator.subscribed_quotes(), vec![instrument.id()]);
2798 assert!(commands.iter().any(|command| matches!(
2799 command,
2800 DataCommand::Subscribe(SubscribeCommand::Quotes(command))
2801 if command.instrument_id == instrument.id()
2802 )));
2803 }
2804
2805 #[rstest]
2806 fn test_submit_order_caches_command(instrument: CryptoPerpetual) {
2807 let (_clock, cache, emulator) = create_emulator();
2808 add_instrument_to_cache(&cache, &instrument);
2809 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2810 let client_order_id = order.client_order_id();
2811 let command = create_submit_order(&instrument, &order);
2812 cache
2813 .borrow_mut()
2814 .add_order(order, None, None, false)
2815 .unwrap();
2816
2817 emulator
2818 .borrow_mut()
2819 .cache_submit_order_command(command.clone());
2820 emulator.borrow_mut().handle_submit_order(&command);
2821
2822 let commands = emulator.borrow().get_submit_order_commands();
2823 assert!(commands.contains_key(&client_order_id));
2824 }
2825
2826 #[rstest]
2827 fn test_handle_submit_order_applies_emulated_event_to_cache(instrument: CryptoPerpetual) {
2828 let (_clock, cache, emulator) = create_emulator();
2829 let risk_events = register_risk_event_handler("RiskEngine.process.emulated");
2830 add_instrument_to_cache(&cache, &instrument);
2831 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2832 let client_order_id = order.client_order_id();
2833 let strategy_id = order.strategy_id();
2834 let command = create_submit_order(&instrument, &order);
2835 cache
2836 .borrow_mut()
2837 .add_order(order, None, None, false)
2838 .unwrap();
2839 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
2840
2841 emulator
2842 .borrow_mut()
2843 .cache_submit_order_command(command.clone());
2844 emulator.borrow_mut().handle_submit_order(&command);
2845 msgbus::unsubscribe_order_events(
2846 format!("events.order.{strategy_id}").into(),
2847 &order_handler,
2848 );
2849 let cache = cache.borrow();
2850 let cached_order = cache.order(&client_order_id).unwrap();
2851 let risk_events = risk_events.get_messages();
2852 let order_events = order_events.borrow();
2853
2854 assert_eq!(cached_order.status(), OrderStatus::Emulated);
2855 assert_eq!(cached_order.event_count(), 2);
2856 assert_eq!(risk_events.len(), 1);
2857 assert!(matches!(risk_events[0], OrderEventAny::Emulated(_)));
2858 assert_eq!(order_events.len(), 1);
2859 assert!(matches!(order_events[0], OrderEventAny::Emulated(_)));
2860 }
2861
2862 #[rstest]
2863 fn test_submit_order_missing_synthetic_trigger_cancels_order(instrument: CryptoPerpetual) {
2864 let (_clock, cache, emulator) = create_emulator();
2865 add_instrument_to_cache(&cache, &instrument);
2866 let synthetic_id = InstrumentId::from("BTC-ETH-INDEX.SYNTH");
2867 let order = OrderTestBuilder::new(OrderType::StopMarket)
2868 .instrument_id(instrument.id())
2869 .side(OrderSide::Buy)
2870 .trigger_price(Price::from("5100.00"))
2871 .quantity(Quantity::from(1))
2872 .emulation_trigger(TriggerType::BidAsk)
2873 .trigger_instrument_id(synthetic_id)
2874 .build();
2875 let client_order_id = order.client_order_id();
2876 let command = create_submit_order(&instrument, &order);
2877 cache
2878 .borrow_mut()
2879 .add_order(order, None, None, false)
2880 .unwrap();
2881
2882 emulator
2883 .borrow_mut()
2884 .cache_submit_order_command(command.clone());
2885 emulator.borrow_mut().handle_submit_order(&command);
2886
2887 let cache = cache.borrow();
2888 let cached_order = cache.order(&client_order_id).unwrap();
2889 let emulator = emulator.borrow();
2890 let commands = emulator.get_submit_order_commands();
2891
2892 assert_eq!(cached_order.status(), OrderStatus::Canceled);
2893 assert!(emulator.get_matching_core(&synthetic_id).is_none());
2894 assert!(emulator.subscribed_quotes().is_empty());
2895 assert!(emulator.subscribed_trades().is_empty());
2896 assert!(!commands.contains_key(&client_order_id));
2897 }
2898
2899 #[rstest]
2900 fn test_submit_order_releases_when_immediately_triggered(instrument: CryptoPerpetual) {
2901 let (_clock, cache, emulator) = create_emulator();
2902 let risk_events =
2903 register_risk_event_handler("RiskEngine.process.submit_immediately_triggered");
2904 add_instrument_to_cache(&cache, &instrument);
2905 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
2906 core.set_bid_raw(Price::from("5099.00"));
2907 core.set_ask_raw(Price::from("5101.00"));
2908 emulator
2909 .borrow_mut()
2910 .matching_cores
2911 .insert(instrument.id(), core);
2912 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2913 let client_order_id = order.client_order_id();
2914 let command = create_submit_order(&instrument, &order);
2915 cache
2916 .borrow_mut()
2917 .add_order(order, None, None, false)
2918 .unwrap();
2919
2920 emulator.borrow_mut().handle_submit_order(&command);
2921
2922 let cache = cache.borrow();
2923 let cached_order = cache.order(&client_order_id).unwrap();
2924 let emulator = emulator.borrow();
2925 let matching_core = emulator.get_matching_core(&instrument.id()).unwrap();
2926 let risk_events = risk_events.get_messages();
2927 assert_eq!(cached_order.status(), OrderStatus::Released);
2928 assert_eq!(cached_order.order_type(), OrderType::Market);
2929 assert!(!matching_core.order_exists(client_order_id));
2930 assert_eq!(emulator.subscribed_strategy_count(), 1);
2931 assert!(
2932 !emulator
2933 .get_submit_order_commands()
2934 .contains_key(&client_order_id)
2935 );
2936 assert_eq!(risk_events.len(), 1);
2937 assert!(matches!(
2938 &risk_events[0],
2939 OrderEventAny::Released(event) if event.client_order_id == client_order_id
2940 ));
2941 }
2942
2943 #[rstest]
2944 fn test_submit_limit_order_releases_when_immediately_fillable(instrument: CryptoPerpetual) {
2945 let (_clock, cache, emulator) = create_emulator();
2946 let risk_events =
2947 register_risk_event_handler("RiskEngine.process.submit_immediately_fillable");
2948 add_instrument_to_cache(&cache, &instrument);
2949 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
2950 core.set_bid_raw(Price::from("5099.00"));
2951 core.set_ask_raw(Price::from("5101.00"));
2952 emulator
2953 .borrow_mut()
2954 .matching_cores
2955 .insert(instrument.id(), core);
2956 let order = OrderTestBuilder::new(OrderType::Limit)
2957 .instrument_id(instrument.id())
2958 .side(OrderSide::Buy)
2959 .price(Price::from("5101.00"))
2960 .quantity(Quantity::from(1))
2961 .emulation_trigger(TriggerType::BidAsk)
2962 .build();
2963 let client_order_id = order.client_order_id();
2964 let command = create_submit_order(&instrument, &order);
2965 cache
2966 .borrow_mut()
2967 .add_order(order, None, None, false)
2968 .unwrap();
2969
2970 emulator.borrow_mut().handle_submit_order(&command);
2971
2972 let cache = cache.borrow();
2973 let cached_order = cache.order(&client_order_id).unwrap();
2974 let emulator = emulator.borrow();
2975 let matching_core = emulator.get_matching_core(&instrument.id()).unwrap();
2976 let risk_events = risk_events.get_messages();
2977 assert_eq!(cached_order.status(), OrderStatus::Released);
2978 assert_eq!(cached_order.order_type(), OrderType::Market);
2979 assert!(!matching_core.order_exists(client_order_id));
2980 assert!(
2981 !emulator
2982 .get_submit_order_commands()
2983 .contains_key(&client_order_id)
2984 );
2985 assert_eq!(risk_events.len(), 1);
2986 assert!(matches!(
2987 &risk_events[0],
2988 OrderEventAny::Released(event) if event.client_order_id == client_order_id
2989 ));
2990 }
2991
2992 #[rstest]
2993 fn test_modify_order_reindexes_matching_state(instrument: CryptoPerpetual) {
2994 let (_clock, cache, emulator) = create_emulator();
2995 add_instrument_to_cache(&cache, &instrument);
2996 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
2997 let client_order_id = order.client_order_id();
2998 let command = create_submit_order(&instrument, &order);
2999 cache
3000 .borrow_mut()
3001 .add_order(order.clone(), None, None, false)
3002 .unwrap();
3003 emulator.borrow_mut().handle_submit_order(&command);
3004 let exec_events =
3005 register_exec_event_handler(cache.clone(), "ExecEngine.process.modify_reindex");
3006 let new_trigger = Price::from("5200.00");
3007 let modify = ModifyOrder::new(
3008 order.trader_id(),
3009 None,
3010 order.strategy_id(),
3011 instrument.id(),
3012 client_order_id,
3013 None,
3014 None,
3015 None,
3016 Some(new_trigger),
3017 UUID4::new(),
3018 0.into(),
3019 None,
3020 None,
3021 );
3022
3023 emulator.borrow_mut().handle_modify_order(&modify);
3024
3025 let cache = cache.borrow();
3026 let cached_order = cache.order(&client_order_id).unwrap();
3027 let emulator = emulator.borrow();
3028 let matching_core = emulator.get_matching_core(&instrument.id()).unwrap();
3029 let match_info = matching_core.get_order(client_order_id).unwrap();
3030 let exec_events = exec_events.get_messages();
3031 assert_eq!(cached_order.trigger_price(), Some(new_trigger));
3032 assert_eq!(match_info.trigger_price, Some(new_trigger));
3033 assert_eq!(matching_core.get_orders().len(), 1);
3034 assert_eq!(exec_events.len(), 1);
3035 assert!(matches!(
3036 &exec_events[0],
3037 OrderEventAny::Updated(event) if event.client_order_id == client_order_id
3038 ));
3039 }
3040
3041 #[rstest]
3042 fn test_modify_order_releases_when_updated_trigger_matches(instrument: CryptoPerpetual) {
3043 let (_clock, cache, emulator) = create_emulator();
3044 let risk_events =
3045 register_risk_event_handler("RiskEngine.process.modify_immediately_triggered");
3046 add_instrument_to_cache(&cache, &instrument);
3047 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3048 core.set_bid_raw(Price::from("5099.00"));
3049 core.set_ask_raw(Price::from("5101.00"));
3050 emulator
3051 .borrow_mut()
3052 .matching_cores
3053 .insert(instrument.id(), core);
3054 let order = OrderTestBuilder::new(OrderType::StopMarket)
3055 .instrument_id(instrument.id())
3056 .side(OrderSide::Buy)
3057 .trigger_price(Price::from("5200.00"))
3058 .quantity(Quantity::from(1))
3059 .emulation_trigger(TriggerType::BidAsk)
3060 .build();
3061 let client_order_id = order.client_order_id();
3062 let command = create_submit_order(&instrument, &order);
3063 cache
3064 .borrow_mut()
3065 .add_order(order.clone(), None, None, false)
3066 .unwrap();
3067 emulator.borrow_mut().handle_submit_order(&command);
3068 risk_events.clear();
3069 let exec_events = register_exec_event_handler(
3070 cache.clone(),
3071 "ExecEngine.process.modify_immediately_triggered",
3072 );
3073 let modify = ModifyOrder::new(
3074 order.trader_id(),
3075 None,
3076 order.strategy_id(),
3077 instrument.id(),
3078 client_order_id,
3079 None,
3080 None,
3081 None,
3082 Some(Price::from("5100.00")),
3083 UUID4::new(),
3084 0.into(),
3085 None,
3086 None,
3087 );
3088
3089 emulator.borrow_mut().handle_modify_order(&modify);
3090
3091 let cache = cache.borrow();
3092 let cached_order = cache.order(&client_order_id).unwrap();
3093 let emulator = emulator.borrow();
3094 let matching_core = emulator.get_matching_core(&instrument.id()).unwrap();
3095 let exec_events = exec_events.get_messages();
3096 let risk_events = risk_events.get_messages();
3097 assert_eq!(cached_order.status(), OrderStatus::Released);
3098 assert_eq!(cached_order.order_type(), OrderType::Market);
3099 assert!(!matching_core.order_exists(client_order_id));
3100 assert!(
3101 !emulator
3102 .get_submit_order_commands()
3103 .contains_key(&client_order_id)
3104 );
3105 assert_eq!(exec_events.len(), 1);
3106 assert!(matches!(
3107 &exec_events[0],
3108 OrderEventAny::Updated(event) if event.client_order_id == client_order_id
3109 ));
3110 assert_eq!(risk_events.len(), 1);
3111 assert!(matches!(
3112 &risk_events[0],
3113 OrderEventAny::Released(event) if event.client_order_id == client_order_id
3114 ));
3115 }
3116
3117 #[rstest]
3118 fn test_update_order_applies_updated_event_to_cache(instrument: CryptoPerpetual) {
3119 let (_clock, cache, emulator) = create_emulator();
3120 let risk_events = register_risk_event_handler("RiskEngine.process.updated");
3121 let mut order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3122 let client_order_id = order.client_order_id();
3123 cache
3124 .borrow_mut()
3125 .add_order(order.clone(), None, None, false)
3126 .unwrap();
3127
3128 emulator
3129 .borrow_mut()
3130 .update_order(&mut order, Quantity::from(2));
3131 let cache = cache.borrow();
3132 let cached_order = cache.order(&client_order_id).unwrap();
3133 let risk_events = risk_events.get_messages();
3134
3135 assert_eq!(order.quantity(), Quantity::from(2));
3136 assert_eq!(cached_order.quantity(), Quantity::from(2));
3137 assert_eq!(cached_order.status(), OrderStatus::Initialized);
3138 assert_eq!(risk_events.len(), 1);
3139 assert!(matches!(risk_events[0], OrderEventAny::Updated(_)));
3140 }
3141
3142 #[rstest]
3143 fn test_fill_market_order_applies_released_event_to_cache(instrument: CryptoPerpetual) {
3144 let (_clock, cache, emulator) = create_emulator();
3145 let risk_events = register_risk_event_handler("RiskEngine.process.released");
3146 add_instrument_to_cache(&cache, &instrument);
3147 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3148 let client_order_id = order.client_order_id();
3149 let strategy_id = order.strategy_id();
3150 let command = create_submit_order(&instrument, &order);
3151 cache
3152 .borrow_mut()
3153 .add_order(order, None, None, false)
3154 .unwrap();
3155
3156 emulator
3157 .borrow_mut()
3158 .cache_submit_order_command(command.clone());
3159 emulator.borrow_mut().handle_submit_order(&command);
3160 risk_events.clear();
3161 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
3162 {
3163 let mut emulator = emulator.borrow_mut();
3164 emulator
3165 .matching_cores
3166 .get_mut(&instrument.id())
3167 .unwrap()
3168 .set_ask_raw(Price::from("5100.00"));
3169 emulator.fill_market_order(client_order_id);
3170 }
3171 msgbus::unsubscribe_order_events(
3172 format!("events.order.{strategy_id}").into(),
3173 &order_handler,
3174 );
3175 let cache = cache.borrow();
3176 let cached_order = cache.order(&client_order_id).unwrap();
3177 let risk_events = risk_events.get_messages();
3178 let order_events = order_events.borrow();
3179
3180 assert_eq!(cached_order.status(), OrderStatus::Released);
3181 assert_eq!(risk_events.len(), 1);
3182 assert!(matches!(risk_events[0], OrderEventAny::Released(_)));
3183 assert_eq!(order_events.len(), 2);
3184 assert!(matches!(order_events[0], OrderEventAny::Initialized(_)));
3185 assert!(matches!(order_events[1], OrderEventAny::Released(_)));
3186 }
3187
3188 #[rstest]
3189 fn test_fill_limit_order_publishes_transformed_initialized_before_released(
3190 instrument: CryptoPerpetual,
3191 ) {
3192 let (_clock, cache, emulator) = create_emulator();
3193 let risk_events = register_risk_event_handler("RiskEngine.process.limit_released");
3194 add_instrument_to_cache(&cache, &instrument);
3195 let order = create_stop_limit_order(&instrument, TriggerType::BidAsk);
3196 let client_order_id = order.client_order_id();
3197 let strategy_id = order.strategy_id();
3198 let command = create_submit_order(&instrument, &order);
3199 cache
3200 .borrow_mut()
3201 .add_order(order, None, None, false)
3202 .unwrap();
3203
3204 emulator
3205 .borrow_mut()
3206 .cache_submit_order_command(command.clone());
3207 emulator.borrow_mut().handle_submit_order(&command);
3208 risk_events.clear();
3209 let (order_handler, order_events) = subscribe_order_topic(strategy_id);
3210 {
3211 let mut emulator = emulator.borrow_mut();
3212 emulator
3213 .matching_cores
3214 .get_mut(&instrument.id())
3215 .unwrap()
3216 .set_ask_raw(Price::from("5100.00"));
3217 emulator.fill_limit_order(client_order_id);
3218 }
3219 msgbus::unsubscribe_order_events(
3220 format!("events.order.{strategy_id}").into(),
3221 &order_handler,
3222 );
3223 let cache = cache.borrow();
3224 let cached_order = cache.order(&client_order_id).unwrap();
3225 let risk_events = risk_events.get_messages();
3226 let order_events = order_events.borrow();
3227
3228 assert_eq!(cached_order.status(), OrderStatus::Released);
3229 assert_eq!(risk_events.len(), 1);
3230 assert!(matches!(risk_events[0], OrderEventAny::Released(_)));
3231 assert_eq!(order_events.len(), 2);
3232 assert!(matches!(order_events[0], OrderEventAny::Initialized(_)));
3233 assert!(matches!(order_events[1], OrderEventAny::Released(_)));
3234 }
3235
3236 #[rstest]
3237 fn test_quote_tick_updates_matching_core_prices(instrument: CryptoPerpetual) {
3238 let (_clock, cache, emulator) = create_emulator();
3239 add_instrument_to_cache(&cache, &instrument);
3240 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3241 let command = create_submit_order(&instrument, &order);
3242 cache
3243 .borrow_mut()
3244 .add_order(order, None, None, false)
3245 .unwrap();
3246 emulator
3247 .borrow_mut()
3248 .cache_submit_order_command(command.clone());
3249 emulator.borrow_mut().handle_submit_order(&command);
3250
3251 let quote = create_quote_tick(&instrument, "5060.00", "5070.00");
3252 emulator.borrow_mut().on_quote_tick(quote);
3253
3254 let core = emulator
3255 .borrow()
3256 .get_matching_core(&instrument.id())
3257 .unwrap();
3258 assert_eq!(core.bid, Some(Price::from("5060.00")));
3259 assert_eq!(core.ask, Some(Price::from("5070.00")));
3260 }
3261
3262 #[rstest]
3263 fn test_trade_tick_updates_matching_core_last_price(instrument: CryptoPerpetual) {
3264 let (_clock, cache, emulator) = create_emulator();
3265 add_instrument_to_cache(&cache, &instrument);
3266 let order = create_stop_market_order(&instrument, TriggerType::LastPrice);
3267 let command = create_submit_order(&instrument, &order);
3268 cache
3269 .borrow_mut()
3270 .add_order(order, None, None, false)
3271 .unwrap();
3272 emulator
3273 .borrow_mut()
3274 .cache_submit_order_command(command.clone());
3275 emulator.borrow_mut().handle_submit_order(&command);
3276
3277 let trade = create_trade_tick(&instrument, "5065.00");
3278 emulator.borrow_mut().on_trade_tick(trade);
3279
3280 let core = emulator
3281 .borrow()
3282 .get_matching_core(&instrument.id())
3283 .unwrap();
3284 assert_eq!(core.last, Some(Price::from("5065.00")));
3285 }
3286
3287 #[rstest]
3288 fn test_cancel_order_removes_from_matching_core(instrument: CryptoPerpetual) {
3289 let (_clock, cache, emulator) = create_emulator();
3290 add_instrument_to_cache(&cache, &instrument);
3291 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3292 let command = create_submit_order(&instrument, &order);
3293 cache
3294 .borrow_mut()
3295 .add_order(order.clone(), None, None, false)
3296 .unwrap();
3297 emulator
3298 .borrow_mut()
3299 .cache_submit_order_command(command.clone());
3300 emulator.borrow_mut().handle_submit_order(&command);
3301
3302 emulator.borrow_mut().cancel_order(&order);
3303
3304 let core = emulator
3305 .borrow()
3306 .get_matching_core(&instrument.id())
3307 .unwrap();
3308 assert!(core.get_orders().is_empty());
3309 }
3310
3311 #[rstest]
3312 fn test_cancel_order_removes_cached_command(instrument: CryptoPerpetual) {
3313 let (_clock, cache, emulator) = create_emulator();
3314 add_instrument_to_cache(&cache, &instrument);
3315 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3316 let client_order_id = order.client_order_id();
3317 let command = create_submit_order(&instrument, &order);
3318 cache
3319 .borrow_mut()
3320 .add_order(order.clone(), None, None, false)
3321 .unwrap();
3322 emulator
3323 .borrow_mut()
3324 .cache_submit_order_command(command.clone());
3325 emulator.borrow_mut().handle_submit_order(&command);
3326
3327 emulator.borrow_mut().cancel_order(&order);
3328
3329 let commands = emulator.borrow().get_submit_order_commands();
3330 assert!(!commands.contains_key(&client_order_id));
3331 }
3332
3333 #[rstest]
3334 fn test_trailing_stop_waits_for_activation_price(instrument: CryptoPerpetual) {
3335 let (_clock, cache, emulator) = create_emulator();
3336 let risk_events = register_risk_event_handler("RiskEngine.process.trailing_activation");
3337 let (handler, exec_commands): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3338 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3339 msgbus::register_trading_command_endpoint(
3340 MessagingSwitchboard::exec_engine_queue_execute(),
3341 handler,
3342 );
3343 add_instrument_to_cache(&cache, &instrument);
3344
3345 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3346 core.set_bid_raw(Price::from("5000.00"));
3347 core.set_ask_raw(Price::from("5001.00"));
3348 emulator
3349 .borrow_mut()
3350 .matching_cores
3351 .insert(instrument.id(), core);
3352
3353 let order = OrderTestBuilder::new(OrderType::TrailingStopMarket)
3354 .instrument_id(instrument.id())
3355 .client_order_id(ClientOrderId::from("O-TRAILING-1"))
3356 .side(OrderSide::Sell)
3357 .quantity(Quantity::from(1))
3358 .activation_price(Price::from("5050.00"))
3359 .trailing_offset(dec!(10))
3360 .trailing_offset_type(TrailingOffsetType::Price)
3361 .emulation_trigger(TriggerType::BidAsk)
3362 .build();
3363 let client_order_id = order.client_order_id();
3364 let command = create_submit_order(&instrument, &order);
3365 cache
3366 .borrow_mut()
3367 .add_order(order, None, None, false)
3368 .unwrap();
3369
3370 emulator
3371 .borrow_mut()
3372 .cache_submit_order_command(command.clone());
3373 emulator.borrow_mut().handle_submit_order(&command);
3374
3375 {
3376 let cache = cache.borrow();
3377 let cached = cache.order(&client_order_id).unwrap();
3378 let is_activated = match &*cached {
3379 OrderAny::TrailingStopMarket(o) => o.is_activated,
3380 _ => panic!("expected trailing stop market order"),
3381 };
3382 assert!(
3383 !is_activated,
3384 "must not activate before the activation price trades"
3385 );
3386 assert!(cached.trigger_price().is_none());
3387 }
3388 assert!(
3389 !risk_events
3390 .get_messages()
3391 .iter()
3392 .any(|e| matches!(e, OrderEventAny::Updated(_))),
3393 "no trailing update before activation"
3394 );
3395
3396 emulator
3398 .borrow_mut()
3399 .on_quote_tick(create_quote_tick(&instrument, "5055.00", "5056.00"));
3400
3401 {
3402 let cache = cache.borrow();
3403 let cached = cache.order(&client_order_id).unwrap();
3404 let is_activated = match &*cached {
3405 OrderAny::TrailingStopMarket(o) => o.is_activated,
3406 _ => panic!("expected trailing stop market order"),
3407 };
3408 assert!(is_activated);
3409 assert_eq!(cached.trigger_price(), Some(Price::from("5045.00")));
3410 }
3411 assert!(exec_commands.get_messages().is_empty());
3412
3413 emulator
3415 .borrow_mut()
3416 .on_quote_tick(create_quote_tick(&instrument, "5044.00", "5045.00"));
3417
3418 let commands = exec_commands.get_messages();
3419 assert_eq!(commands.len(), 1);
3420 assert!(matches!(
3421 &commands[0],
3422 TradingCommand::SubmitOrder(command) if command.client_order_id == client_order_id
3423 ));
3424 }
3425
3426 #[rstest]
3427 fn test_pending_activation_trailing_stop_limit_not_matched_as_plain_limit(
3428 instrument: CryptoPerpetual,
3429 ) {
3430 let (_clock, cache, emulator) = create_emulator();
3431 let (handler, exec_commands): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3432 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3433 msgbus::register_trading_command_endpoint(
3434 MessagingSwitchboard::exec_engine_queue_execute(),
3435 handler,
3436 );
3437 add_instrument_to_cache(&cache, &instrument);
3438
3439 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3441 core.set_bid_raw(Price::from("5000.00"));
3442 core.set_ask_raw(Price::from("5001.00"));
3443 emulator
3444 .borrow_mut()
3445 .matching_cores
3446 .insert(instrument.id(), core);
3447
3448 let order = OrderTestBuilder::new(OrderType::TrailingStopLimit)
3449 .instrument_id(instrument.id())
3450 .client_order_id(ClientOrderId::from("O-TRAILING-LIMIT-1"))
3451 .side(OrderSide::Sell)
3452 .price(Price::from("5010.00"))
3453 .quantity(Quantity::from(1))
3454 .activation_price(Price::from("5050.00"))
3455 .trailing_offset(dec!(10))
3456 .trailing_offset_type(TrailingOffsetType::Price)
3457 .limit_offset(dec!(5))
3458 .emulation_trigger(TriggerType::BidAsk)
3459 .build();
3460 let client_order_id = order.client_order_id();
3461 let command = create_submit_order(&instrument, &order);
3462 cache
3463 .borrow_mut()
3464 .add_order(order, None, None, false)
3465 .unwrap();
3466
3467 emulator
3468 .borrow_mut()
3469 .cache_submit_order_command(command.clone());
3470 emulator.borrow_mut().handle_submit_order(&command);
3471
3472 emulator
3475 .borrow_mut()
3476 .on_quote_tick(create_quote_tick(&instrument, "4985.00", "4986.00"));
3477
3478 assert!(exec_commands.get_messages().is_empty());
3479 assert!(
3480 emulator
3481 .borrow()
3482 .get_submit_order_commands()
3483 .contains_key(&client_order_id),
3484 "order must remain held by the emulator before activation"
3485 );
3486 }
3487
3488 #[rstest]
3489 fn test_trailing_stop_with_preset_trigger_activates_and_triggers(instrument: CryptoPerpetual) {
3490 let (_clock, cache, emulator) = create_emulator();
3491 let (handler, exec_commands): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3492 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3493 msgbus::register_trading_command_endpoint(
3494 MessagingSwitchboard::exec_engine_queue_execute(),
3495 handler,
3496 );
3497 add_instrument_to_cache(&cache, &instrument);
3498
3499 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3500 core.set_bid_raw(Price::from("5000.00"));
3501 core.set_ask_raw(Price::from("5001.00"));
3502 emulator
3503 .borrow_mut()
3504 .matching_cores
3505 .insert(instrument.id(), core);
3506
3507 let order = OrderTestBuilder::new(OrderType::TrailingStopMarket)
3509 .instrument_id(instrument.id())
3510 .client_order_id(ClientOrderId::from("O-TRAILING-2"))
3511 .side(OrderSide::Sell)
3512 .quantity(Quantity::from(1))
3513 .trigger_price(Price::from("5045.00"))
3514 .activation_price(Price::from("5050.00"))
3515 .trailing_offset(dec!(10))
3516 .trailing_offset_type(TrailingOffsetType::Price)
3517 .emulation_trigger(TriggerType::BidAsk)
3518 .build();
3519 let client_order_id = order.client_order_id();
3520 let command = create_submit_order(&instrument, &order);
3521 cache
3522 .borrow_mut()
3523 .add_order(order, None, None, false)
3524 .unwrap();
3525
3526 emulator
3527 .borrow_mut()
3528 .cache_submit_order_command(command.clone());
3529 emulator.borrow_mut().handle_submit_order(&command);
3530
3531 assert!(
3532 exec_commands.get_messages().is_empty(),
3533 "preset trigger must not release before activation"
3534 );
3535
3536 emulator
3539 .borrow_mut()
3540 .on_quote_tick(create_quote_tick(&instrument, "5055.00", "5056.00"));
3541
3542 assert!(exec_commands.get_messages().is_empty());
3543
3544 emulator
3546 .borrow_mut()
3547 .on_quote_tick(create_quote_tick(&instrument, "5044.00", "5045.00"));
3548
3549 let commands = exec_commands.get_messages();
3550 assert_eq!(commands.len(), 1);
3551 assert!(matches!(
3552 &commands[0],
3553 TradingCommand::SubmitOrder(command) if command.client_order_id == client_order_id
3554 ));
3555 }
3556
3557 #[rstest]
3558 fn test_trailing_stop_activates_despite_trailing_calculate_error(instrument: CryptoPerpetual) {
3559 let (_clock, cache, emulator) = create_emulator();
3560 let (handler, exec_commands): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3561 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3562 msgbus::register_trading_command_endpoint(
3563 MessagingSwitchboard::exec_engine_queue_execute(),
3564 handler,
3565 );
3566 add_instrument_to_cache(&cache, &instrument);
3567
3568 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3570 core.set_bid_raw(Price::from("5000.00"));
3571 core.set_ask_raw(Price::from("5001.00"));
3572 emulator
3573 .borrow_mut()
3574 .matching_cores
3575 .insert(instrument.id(), core);
3576
3577 let order = OrderTestBuilder::new(OrderType::TrailingStopMarket)
3578 .instrument_id(instrument.id())
3579 .client_order_id(ClientOrderId::from("O-TRAILING-3"))
3580 .side(OrderSide::Sell)
3581 .quantity(Quantity::from(1))
3582 .trigger_price(Price::from("5045.00"))
3583 .trigger_type(TriggerType::LastOrBidAsk)
3584 .activation_price(Price::from("5050.00"))
3585 .trailing_offset(dec!(10))
3586 .trailing_offset_type(TrailingOffsetType::Price)
3587 .emulation_trigger(TriggerType::BidAsk)
3588 .build();
3589 let client_order_id = order.client_order_id();
3590 let command = create_submit_order(&instrument, &order);
3591 cache
3592 .borrow_mut()
3593 .add_order(order, None, None, false)
3594 .unwrap();
3595
3596 emulator
3597 .borrow_mut()
3598 .cache_submit_order_command(command.clone());
3599 emulator.borrow_mut().handle_submit_order(&command);
3600
3601 emulator
3604 .borrow_mut()
3605 .on_quote_tick(create_quote_tick(&instrument, "5055.00", "5056.00"));
3606 emulator
3607 .borrow_mut()
3608 .on_quote_tick(create_quote_tick(&instrument, "5044.00", "5045.00"));
3609
3610 let commands = exec_commands.get_messages();
3611 assert_eq!(commands.len(), 1);
3612 assert!(matches!(
3613 &commands[0],
3614 TradingCommand::SubmitOrder(command) if command.client_order_id == client_order_id
3615 ));
3616 }
3617
3618 #[rstest]
3619 fn test_cancel_emulated_oco_leg_cancels_sibling(instrument: CryptoPerpetual) {
3620 let (_clock, cache, emulator) = create_emulator();
3621 add_instrument_to_cache(&cache, &instrument);
3622 let client_order_id_a = ClientOrderId::from("O-OCO-A");
3623 let client_order_id_b = ClientOrderId::from("O-OCO-B");
3624 let order_a = OrderTestBuilder::new(OrderType::StopMarket)
3625 .instrument_id(instrument.id())
3626 .strategy_id(StrategyId::from("STRATEGY-001"))
3627 .client_order_id(client_order_id_a)
3628 .side(OrderSide::Sell)
3629 .trigger_price(Price::from("4900.00"))
3630 .quantity(Quantity::from(1))
3631 .emulation_trigger(TriggerType::BidAsk)
3632 .contingency_type(ContingencyType::Oco)
3633 .linked_order_ids(vec![client_order_id_b])
3634 .build();
3635 let order_b = OrderTestBuilder::new(OrderType::StopMarket)
3636 .instrument_id(instrument.id())
3637 .strategy_id(StrategyId::from("STRATEGY-001"))
3638 .client_order_id(client_order_id_b)
3639 .side(OrderSide::Sell)
3640 .trigger_price(Price::from("4950.00"))
3641 .quantity(Quantity::from(1))
3642 .emulation_trigger(TriggerType::BidAsk)
3643 .contingency_type(ContingencyType::Oco)
3644 .linked_order_ids(vec![client_order_id_a])
3645 .build();
3646 cache
3647 .borrow_mut()
3648 .add_order(order_a.clone(), None, None, false)
3649 .unwrap();
3650 cache
3651 .borrow_mut()
3652 .add_order(order_b.clone(), None, None, false)
3653 .unwrap();
3654
3655 for order in [&order_a, &order_b] {
3656 let command = create_submit_order(&instrument, order);
3657 emulator
3658 .borrow_mut()
3659 .cache_submit_order_command(command.clone());
3660 emulator.borrow_mut().handle_submit_order(&command);
3661 }
3662 assert_eq!(
3663 cache.borrow().order(&client_order_id_a).unwrap().status(),
3664 OrderStatus::Emulated
3665 );
3666 assert_eq!(
3667 cache.borrow().order(&client_order_id_b).unwrap().status(),
3668 OrderStatus::Emulated
3669 );
3670
3671 let cancel = CancelOrder::new(
3672 order_a.trader_id(),
3673 None,
3674 order_a.strategy_id(),
3675 instrument.id(),
3676 client_order_id_a,
3677 None,
3678 UUID4::new(),
3679 0.into(),
3680 None,
3681 None,
3682 );
3683 msgbus::send_trading_command(
3684 MessagingSwitchboard::order_emulator_execute(),
3685 TradingCommand::CancelOrder(cancel),
3686 );
3687
3688 let cache = cache.borrow();
3689 assert_eq!(
3690 cache.order(&client_order_id_a).unwrap().status(),
3691 OrderStatus::Canceled
3692 );
3693 assert_eq!(
3694 cache.order(&client_order_id_b).unwrap().status(),
3695 OrderStatus::Canceled,
3696 "OCO sibling must be canceled through the deferred self-published event"
3697 );
3698 }
3699
3700 #[rstest]
3701 fn test_reentrant_command_from_event_handler_does_not_panic(instrument: CryptoPerpetual) {
3702 let (_clock, cache, emulator) = create_emulator();
3703 add_instrument_to_cache(&cache, &instrument);
3704 let order = create_stop_market_order(&instrument, TriggerType::BidAsk);
3705 let client_order_id = order.client_order_id();
3706 let strategy_id = order.strategy_id();
3707 let command = create_submit_order(&instrument, &order);
3708 cache
3709 .borrow_mut()
3710 .add_order(order.clone(), None, None, false)
3711 .unwrap();
3712 emulator
3713 .borrow_mut()
3714 .cache_submit_order_command(command.clone());
3715
3716 let cancel = CancelOrder::new(
3719 order.trader_id(),
3720 None,
3721 strategy_id,
3722 instrument.id(),
3723 client_order_id,
3724 None,
3725 UUID4::new(),
3726 0.into(),
3727 None,
3728 None,
3729 );
3730 let handler = TypedHandler::from(move |event: &OrderEventAny| {
3731 if matches!(event, OrderEventAny::Emulated(_)) {
3732 msgbus::send_trading_command(
3733 MessagingSwitchboard::order_emulator_execute(),
3734 TradingCommand::CancelOrder(cancel.clone()),
3735 );
3736 }
3737 });
3738 msgbus::subscribe_order_events(format!("events.order.{strategy_id}").into(), handler, None);
3739
3740 msgbus::send_trading_command(
3742 MessagingSwitchboard::order_emulator_execute(),
3743 TradingCommand::SubmitOrder(command),
3744 );
3745
3746 assert!(
3747 cache.borrow().order(&client_order_id).unwrap().is_closed(),
3748 "deferred cancel must be processed after the active call completes"
3749 );
3750 }
3751
3752 #[rstest]
3753 fn test_released_order_preserves_chronological_event_history(instrument: CryptoPerpetual) {
3754 let (_clock, cache, emulator) = create_emulator();
3755 add_instrument_to_cache(&cache, &instrument);
3756 let mut core = OrderMatchingCore::new(instrument.id(), instrument.price_increment());
3757 core.set_bid_raw(Price::from("5099.00"));
3758 core.set_ask_raw(Price::from("5101.00"));
3759 emulator
3760 .borrow_mut()
3761 .matching_cores
3762 .insert(instrument.id(), core);
3763
3764 let order = create_stop_limit_order(&instrument, TriggerType::BidAsk);
3765 let client_order_id = order.client_order_id();
3766 let original_init = OrderEventAny::Initialized(order.init_event().clone());
3767 let command = create_submit_order(&instrument, &order);
3768 cache
3769 .borrow_mut()
3770 .add_order(order, None, None, false)
3771 .unwrap();
3772
3773 emulator
3774 .borrow_mut()
3775 .cache_submit_order_command(command.clone());
3776 emulator.borrow_mut().handle_submit_order(&command);
3777
3778 let cache = cache.borrow();
3779 let transformed = cache.order(&client_order_id).unwrap();
3780 let events = transformed.events();
3781
3782 assert_eq!(
3783 events.first(),
3784 Some(&&original_init),
3785 "released order history must start with the original OrderInitialized event, was {events:?}",
3786 );
3787 }
3788
3789 #[rstest]
3790 fn test_on_start_cancels_emulated_child_of_closed_positionless_parent(
3791 instrument: CryptoPerpetual,
3792 ) {
3793 let (_clock, cache, emulator) = create_emulator();
3794 add_instrument_to_cache(&cache, &instrument);
3795
3796 let mut parent = OrderTestBuilder::new(OrderType::Limit)
3797 .instrument_id(instrument.id())
3798 .client_order_id(ClientOrderId::from("O-PARENT"))
3799 .side(OrderSide::Buy)
3800 .price(Price::from("5000.00"))
3801 .quantity(Quantity::from(1))
3802 .submit(true)
3803 .build();
3804 parent
3805 .apply(OrderEventAny::Canceled(OrderCanceled::new(
3806 parent.trader_id(),
3807 parent.strategy_id(),
3808 parent.instrument_id(),
3809 parent.client_order_id(),
3810 UUID4::new(),
3811 0.into(),
3812 0.into(),
3813 false,
3814 None,
3815 None,
3816 None,
3817 )))
3818 .unwrap();
3819 assert!(parent.is_closed());
3820 assert!(parent.position_id().is_none());
3821
3822 let mut child = OrderTestBuilder::new(OrderType::StopMarket)
3823 .instrument_id(instrument.id())
3824 .client_order_id(ClientOrderId::from("O-CHILD"))
3825 .side(OrderSide::Sell)
3826 .trigger_price(Price::from("4900.00"))
3827 .quantity(Quantity::from(1))
3828 .emulation_trigger(TriggerType::BidAsk)
3829 .parent_order_id(parent.client_order_id())
3830 .build();
3831 child
3832 .apply(OrderEventAny::Emulated(OrderEmulated::new(
3833 child.trader_id(),
3834 child.strategy_id(),
3835 child.instrument_id(),
3836 child.client_order_id(),
3837 UUID4::new(),
3838 0.into(),
3839 0.into(),
3840 )))
3841 .unwrap();
3842 assert_eq!(child.status(), OrderStatus::Emulated);
3843
3844 cache
3845 .borrow_mut()
3846 .add_order(parent.clone(), None, None, false)
3847 .unwrap();
3848 cache
3849 .borrow_mut()
3850 .add_order(child.clone(), None, None, false)
3851 .unwrap();
3852
3853 emulator.borrow_mut().on_start().unwrap();
3854
3855 assert!(
3856 cache
3857 .borrow()
3858 .order(&child.client_order_id())
3859 .unwrap()
3860 .is_closed(),
3861 "emulated child of a closed position-less parent must be canceled on reactivation"
3862 );
3863 }
3864}