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