Skip to main content

nautilus_execution/order_emulator/
emulator.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16use 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    /// Caches a submit order command for emulation tracking.
162    pub fn cache_submit_order_command(&mut self, command: SubmitOrder) {
163        self.manager.cache_submit_order_command(command);
164    }
165
166    /// Subscribes to quote data for the given instrument.
167    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    /// Subscribes to trade data for the given instrument.
197    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    /// Reactivates emulated orders from cache on start.
287    ///
288    /// # Errors
289    ///
290    /// Returns an error if processing fails.
291    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; // No longer emulated
311            }
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; // Parent already closed
328                }
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; // Process contingency order later once parent triggered
335                }
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, // params
359                UUID4::new(),
360                self.clock.borrow().timestamp_ns(),
361                None, // correlation_id
362            );
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        // else: Order not in cache yet
387
388        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        // Sort so the emitted unsubscribe commands, and the UUID4 draws they consume,
419        // follow the same sequence across runs; the drained sets are AHash-backed.
420        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    /// # Panics
541    ///
542    /// Panics if the emulation trigger type is not set or if the order is not in the cache.
543    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        // Get or create matching core
581        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            // Handle synthetic instruments
591            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        // Update trailing stop
630        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            // Not yet activated: held inert until the activation price is touched
645        }
646
647        self.check_monitoring(strategy_id, position_id);
648
649        // Check if immediately marketable
650        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        // Handle data subscriptions
669        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        // Check if order was already released
691        if !self
692            .manager
693            .get_submit_order_commands()
694            .contains_key(&order.client_order_id())
695        {
696            return; // Already released
697        }
698
699        // Hold in matching core
700        matching_core.add_order(match_info);
701
702        // Generate emulated event if needed
703        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        // Since we are cloning the matching core, we need to insert it back into the original hashmap
733        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; // Process contingency order later once parent triggered
759                }
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            // A synchronous event handler may supersede this update before dispatch returns
815            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            // Order not held in the emulator
909            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 = &quote.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        // Process bid actions before ask actions so cross-side
1064        // contingencies (OCO/OUO) mutate state between sides
1065        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        // Re-snapshot orders after actions to avoid stale trailing stop updates
1087        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            // Subscribe to strategy order events
1179            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    /// Validates market data availability for order release.
1198    ///
1199    /// Returns `Some(released_price)` if market data is available, `None` otherwise.
1200    /// Logs appropriate warnings when market data is not yet available.
1201    ///
1202    /// Does NOT pop the submit order command - caller must do that and handle missing command
1203    /// according to their contract (panic for market orders, return for limit orders).
1204    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    /// # Panics
1227    ///
1228    /// Panics if the order type is invalid for a stop order.
1229    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    /// # Panics
1258    ///
1259    /// Panics if a limit order has no price.
1260    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; // Order stays queued for retry
1290            }
1291        };
1292
1293        let released_price =
1294            match self.validate_release(&order, matching_core, trigger_instrument_id) {
1295                Some(price) => price,
1296                None => return, // Order stays queued for retry
1297            };
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, // Order already released
1305        };
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            // Transform order
1313            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            // TODO: fix
1348            // let triggered_price = order.trigger_price();
1349            // if triggered_price.is_some() {
1350            //     transformed.trigger_price() = (triggered_price.unwrap());
1351            // }
1352
1353            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            // Publish event
1402            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    /// # Panics
1416    ///
1417    /// Panics if a market order command is missing.
1418    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; // Order stays queued for retry
1443            }
1444        };
1445
1446        let released_price =
1447            match self.validate_release(&order, matching_core, trigger_instrument_id) {
1448                Some(price) => price,
1449                None => return, // Order stays queued for retry
1450            };
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            // Transform order
1465            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            // Publish event
1533            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; // Not yet activated
1578        }
1579
1580        if !was_activated {
1581            // Activation transitioned: re-key the held entry from its inert
1582            // pre-activation state regardless of the trailing outcome below.
1583            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            // A trailing stop with no trigger price yet must stay inert:
1649            // its limit price must never match as a plain limit order.
1650            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    /// Returns true if the trailing stop is (or has just become) activated.
1668    ///
1669    /// With no `activation_price` the order activates immediately at the current
1670    /// market price; otherwise it activates only once the market touches the
1671    /// activation price.
1672    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, // correlation_id
1941        )
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()), &quote);
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, &quote_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(&quote_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()), &quote);
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        // Registered in an order that is neither sorted nor reverse sorted, and five
2687        // entries leave a 1-in-120 chance that AHash order coincides with the expectation.
2688        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        // Market rallies through the 5050 activation price: activates and starts trailing
3397        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        // Market falls through the trailed trigger: order releases
3414        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        // Market is already through the limit price, but below the activation price
3440        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        // The market falls below where the trailed trigger would sit (4990)
3473        // while staying below the 5050 activation price: must not release.
3474        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        // Both prices preset: the 5045 trigger must stay inert until 5050 activates
3508        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        // Rally touches the 5050 activation price; the trailed candidate does
3537        // not improve on the preset trigger, so no price update is emitted.
3538        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        // Bid falls through the preset trigger: the order must release
3545        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        // Quote-only data: `LastOrBidAsk` trailing calculation errors without a last
3569        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        // Activation touches while the trailing calculation errors: the preset
3602        // trigger must still be honored once the market crosses it.
3603        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        // A strategy-style subscriber issuing an emulator command synchronously
3717        // from within its own order-event handler.
3718        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        // Must not panic on the reentrant borrow; the command is deferred and drained
3741        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}