Skip to main content

nautilus_trading/strategy/
mod.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
16pub mod api;
17pub mod config;
18pub mod core;
19
20pub use core::{StrategyCore, StrategyNative};
21use std::panic::{AssertUnwindSafe, catch_unwind};
22
23use ahash::AHashSet;
24pub use api::{OrderApi, PortfolioApi};
25pub use config::{ImportableStrategyConfig, StrategyConfig};
26use nautilus_common::{
27    actor::DataActor,
28    component::Component,
29    enums::ComponentState,
30    logging::{CMD, EVT, RECV, SEND},
31    messages::execution::{
32        BatchCancelOrders, BatchModifyOrders, CancelAllOrders, CancelOrder, ModifyOrder,
33        QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList, TradingCommand,
34    },
35    msgbus::{self, MessagingSwitchboard},
36    timer::TimeEvent,
37};
38use nautilus_core::{Params, UUID4};
39use nautilus_execution::order_manager::OrderManagerAction;
40use nautilus_model::{
41    enums::{OrderSide, OrderStatus, PositionSide, TimeInForce},
42    events::{
43        OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied, OrderEmulated,
44        OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled, OrderInitialized,
45        OrderModifyRejected, OrderPendingCancel, OrderPendingUpdate, OrderRejected, OrderReleased,
46        OrderSubmitted, OrderTriggered, OrderUpdated, PositionChanged, PositionClosed,
47        PositionEvent, PositionOpened,
48    },
49    identifiers::{
50        AccountId, ClientId, ClientOrderId, ExecAlgorithmId, InstrumentId, PositionId, StrategyId,
51        TraderId,
52    },
53    orders::{
54        LIMIT_ORDER_TYPES, Order, OrderAny, OrderCore, OrderError, OrderList, STOP_ORDER_TYPES,
55    },
56    position::Position,
57    types::{Price, Quantity},
58};
59use ustr::Ustr;
60
61/// Describes one child update in a batch modify request.
62pub type BatchModifyOrder = (
63    ClientOrderId,
64    Option<Quantity>,
65    Option<Price>,
66    Option<Price>,
67);
68
69/// Core trait for implementing trading strategies in NautilusTrader.
70///
71/// Strategies are specialized [`DataActor`]s that combine data ingestion capabilities with
72/// order and position management functionality. By implementing this trait,
73/// custom strategies gain access to the full trading execution stack including order
74/// submission, modification, cancellation, and position management.
75///
76/// # Key Capabilities
77///
78/// - All [`DataActor`] capabilities (data subscriptions, event handling, timers).
79/// - Order lifecycle management (submit, modify, cancel).
80/// - Position management (open, close, monitor).
81/// - Access to the trading cache and portfolio.
82/// - Event routing for orders and emulator events.
83///
84/// # Implementation
85///
86/// Use the `nautilus_strategy!` macro to generate the native runtime wiring
87/// and `Strategy` implementations. Normal strategy logic should call facade
88/// methods such as `strategy_id()`, `clock()`, `cache()`, `order()`, and
89/// `portfolio()`. Native runtime code that needs the internal core should use
90/// [`StrategyNative`]. For strategies that override additional trait methods,
91/// pass them in a block:
92///
93/// ```ignore
94/// nautilus_strategy!(MyStrategy, {
95///     fn on_order_rejected(&mut self, event: OrderRejected) {
96///         // custom handling
97///     }
98/// });
99/// ```
100///
101/// Default methods that read or mutate native runtime state carry explicit
102/// [`StrategyNative`] and [`Component`] bounds. Implementations that only need
103/// behavioral callbacks do not own or implement native runtime state.
104pub trait Strategy: DataActor {
105    /// Returns the external order claims for this strategy.
106    ///
107    /// These are instrument IDs whose external orders should be claimed by this strategy
108    /// during reconciliation.
109    fn external_order_claims(&self) -> Option<Vec<InstrumentId>> {
110        None
111    }
112
113    /// Returns the runtime strategy ID, when configured or registered.
114    fn strategy_id(&self) -> Option<StrategyId>
115    where
116        Self: StrategyNative,
117    {
118        StrategyNative::strategy_core(self).strategy_id()
119    }
120
121    /// Returns the user-facing order creation API.
122    fn order(&self) -> OrderApi<'_>
123    where
124        Self: StrategyNative,
125    {
126        StrategyNative::strategy_core(self).order()
127    }
128
129    /// Returns the user-facing portfolio read API.
130    fn portfolio(&self) -> PortfolioApi<'_>
131    where
132        Self: StrategyNative,
133    {
134        StrategyNative::strategy_core(self).portfolio_api()
135    }
136
137    /// Submits an order.
138    ///
139    /// # Errors
140    ///
141    /// Returns an error if the strategy is not registered or order submission fails.
142    fn submit_order(
143        &mut self,
144        order: OrderAny,
145        position_id: Option<PositionId>,
146        client_id: Option<ClientId>,
147        params: Option<Params>,
148    ) -> anyhow::Result<()>
149    where
150        Self: StrategyNative,
151    {
152        let core = StrategyNative::strategy_core_mut(self);
153
154        let trader_id = registered_trader_id(core)?;
155        let strategy_id = registered_strategy_id(core)?;
156        let ts_init = core.clock_mut().timestamp_ns();
157
158        if order.status() != OrderStatus::Initialized {
159            anyhow::bail!(
160                "Order denied: invalid status for {}, expected INITIALIZED",
161                order.client_order_id()
162            );
163        }
164
165        let market_exit_tag = core.market_exit_tag;
166        let is_market_exit_order = order
167            .tags()
168            .is_some_and(|tags| tags.contains(&market_exit_tag));
169        let should_deny_for_market_exit =
170            core.is_exiting && !order.is_reduce_only() && !is_market_exit_order;
171
172        if should_deny_for_market_exit {
173            self.deny_order(&order, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
174            return Ok(());
175        }
176
177        let core = StrategyNative::strategy_core_mut(self);
178        let params = params.filter(|params| !params.is_empty());
179
180        {
181            let cache_rc = core.cache_rc();
182            let mut cache = cache_rc.try_borrow_mut().map_err(|_| {
183                anyhow::anyhow!(
184                    "Cannot submit order {}: cache is currently borrowed",
185                    order.client_order_id()
186                )
187            })?;
188            cache.add_order(order.clone(), position_id, client_id, true)?;
189        }
190
191        publish_order_initialized(&order);
192
193        let command = SubmitOrder::new(
194            trader_id,
195            client_id,
196            strategy_id,
197            order.instrument_id(),
198            order.client_order_id(),
199            order.init_event().clone(),
200            order.exec_algorithm_id(),
201            position_id,
202            params,
203            UUID4::new(),
204            ts_init,
205            None, // correlation_id
206        );
207
208        if order.emulation_trigger().is_some() {
209            send_emulator_command(TradingCommand::SubmitOrder(command));
210        } else if let Some(exec_algorithm_id) = order.exec_algorithm_id() {
211            send_algo_command(command, exec_algorithm_id);
212        } else {
213            send_risk_command(TradingCommand::SubmitOrder(command));
214        }
215
216        self.set_gtd_expiry(&order)?;
217        Ok(())
218    }
219
220    /// Submits an order list.
221    ///
222    /// # Errors
223    ///
224    /// Returns an error if the strategy is not registered, the order list is invalid,
225    /// or order list submission fails.
226    fn submit_order_list(
227        &mut self,
228        mut orders: Vec<OrderAny>,
229        position_id: Option<PositionId>,
230        client_id: Option<ClientId>,
231        params: Option<Params>,
232    ) -> anyhow::Result<()>
233    where
234        Self: StrategyNative,
235    {
236        if orders.is_empty() {
237            log::error!("OrderList denied: no orders to submit");
238            anyhow::bail!("OrderList denied: no orders to submit");
239        }
240
241        for order in &orders {
242            if order.status() != OrderStatus::Initialized {
243                anyhow::bail!(
244                    "Order in list denied: invalid status for {}, expected INITIALIZED",
245                    order.client_order_id()
246                );
247            }
248        }
249
250        let first_venue = orders[0].instrument_id().venue;
251        for order in &orders {
252            if order.instrument_id().venue != first_venue {
253                anyhow::bail!(
254                    "OrderList denied: orders must share the same venue; \
255                     expected {first_venue}, found {} on {}",
256                    order.instrument_id().venue,
257                    order.client_order_id(),
258                );
259            }
260        }
261
262        let should_deny = {
263            let core = StrategyNative::strategy_core_mut(self);
264            let tag = core.market_exit_tag;
265            core.is_exiting
266                && orders.iter().any(|o| {
267                    !o.is_reduce_only() && !o.tags().is_some_and(|tags| tags.contains(&tag))
268                })
269        };
270
271        if should_deny {
272            self.deny_order_list(&orders, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
273            return Ok(());
274        }
275
276        let core = StrategyNative::strategy_core_mut(self);
277
278        let trader_id = registered_trader_id(core)?;
279        let strategy_id = registered_strategy_id(core)?;
280        let ts_init = core.clock_mut().timestamp_ns();
281
282        // TODO: Replace with fluent builder API for order list construction
283        let order_list = if orders.first().is_some_and(|o| o.order_list_id().is_some()) {
284            OrderList::from_orders(&orders, ts_init)
285        } else {
286            core.order_factory().create_list(&mut orders, ts_init)
287        };
288
289        if let Err(e) = order_list.validate() {
290            log::error!("OrderList denied: {e}");
291            anyhow::bail!("OrderList denied: {e}");
292        }
293
294        {
295            let cache_rc = core.cache_rc();
296            let mut cache = cache_rc.try_borrow_mut().map_err(|_| {
297                anyhow::anyhow!(
298                    "Cannot submit order list {}: cache is currently borrowed",
299                    order_list.id
300                )
301            })?;
302
303            if cache.order_list_exists(&order_list.id) {
304                anyhow::bail!("OrderList denied: duplicate {}", order_list.id);
305            }
306
307            for order in &orders {
308                if cache.order_exists(&order.client_order_id()) {
309                    anyhow::bail!(
310                        "Order in list denied: duplicate {}",
311                        order.client_order_id()
312                    );
313                }
314            }
315
316            cache.add_order_list(order_list.clone())?;
317            for order in &orders {
318                cache.add_order(order.clone(), position_id, client_id, true)?;
319            }
320        }
321
322        for order in &orders {
323            publish_order_initialized(order);
324        }
325
326        let params = params.filter(|params| !params.is_empty());
327
328        let first_order = orders.first();
329        let order_inits: Vec<_> = orders.iter().map(|o| o.init_event().clone()).collect();
330        let exec_algorithm_id = first_order.and_then(Order::exec_algorithm_id);
331
332        let command = SubmitOrderList::new(
333            trader_id,
334            client_id,
335            strategy_id,
336            order_list,
337            order_inits,
338            exec_algorithm_id,
339            position_id,
340            params,
341            UUID4::new(),
342            ts_init,
343            None, // correlation_id
344        );
345
346        let has_emulated_order = orders
347            .iter()
348            .any(|o| o.emulation_trigger().is_some() || o.is_emulated());
349
350        if has_emulated_order {
351            send_emulator_command(TradingCommand::SubmitOrderList(command));
352        } else if let Some(algo_id) = exec_algorithm_id {
353            let endpoint = format!("{algo_id}.execute");
354            msgbus::send_any(endpoint.into(), &TradingCommand::SubmitOrderList(command));
355        } else {
356            send_risk_command(TradingCommand::SubmitOrderList(command));
357        }
358
359        for order in &orders {
360            self.set_gtd_expiry(order)?;
361        }
362
363        Ok(())
364    }
365
366    /// Modifies an order.
367    ///
368    /// # Errors
369    ///
370    /// Returns an error if the strategy is not registered or order modification fails.
371    fn modify_order(
372        &mut self,
373        client_order_id: ClientOrderId,
374        quantity: Option<Quantity>,
375        price: Option<Price>,
376        trigger_price: Option<Price>,
377        client_id: Option<ClientId>,
378        params: Option<Params>,
379    ) -> anyhow::Result<()>
380    where
381        Self: StrategyNative,
382    {
383        let (trader_id, strategy_id) = {
384            let core = StrategyNative::strategy_core_mut(self);
385            (registered_trader_id(core)?, registered_strategy_id(core)?)
386        };
387
388        let params = params.filter(|params| !params.is_empty());
389
390        // TODO: Snapshot the order from the cache. See `cancel_order` for the rationale.
391        let order = StrategyNative::strategy_core_mut(self)
392            .cache_rc()
393            .borrow()
394            .try_order_owned(&client_order_id)
395            .map_err(|e| anyhow::anyhow!("Cannot modify order: {e}"))?;
396
397        let mut updating = false;
398
399        if quantity.is_some_and(|q| q != order.quantity() || order.is_pending_update()) {
400            updating = true;
401        }
402
403        if let Some(price) = price {
404            if !LIMIT_ORDER_TYPES.contains(&order.order_type()) {
405                anyhow::bail!("{} orders do not have a LIMIT price", order.order_type());
406            }
407
408            if Some(price) != order.price() {
409                updating = true;
410            }
411        }
412
413        if let Some(trigger_price) = trigger_price {
414            if !STOP_ORDER_TYPES.contains(&order.order_type()) {
415                anyhow::bail!(
416                    "{} orders do not have a STOP trigger price",
417                    order.order_type()
418                );
419            }
420
421            if Some(trigger_price) != order.trigger_price() {
422                updating = true;
423            }
424        }
425
426        if !updating {
427            log::error!(
428                "Cannot create command ModifyOrder: quantity, price, and trigger were either None \
429                or the same as existing values"
430            );
431            return Ok(());
432        }
433
434        if order.is_closed() || order.is_pending_cancel() {
435            log::warn!(
436                "Cannot create command ModifyOrder: state is {:?}, {order:?}",
437                order.status()
438            );
439            return Ok(());
440        }
441
442        if !self.mark_order_pending_update(&order)? {
443            return Ok(());
444        }
445
446        let command = ModifyOrder::new(
447            trader_id,
448            client_id,
449            strategy_id,
450            order.instrument_id(),
451            order.client_order_id(),
452            order.venue_order_id(),
453            quantity,
454            price,
455            trigger_price,
456            UUID4::new(),
457            StrategyNative::strategy_core_mut(self)
458                .clock_mut()
459                .timestamp_ns(),
460            params,
461            None, // correlation_id
462        );
463
464        if order.emulation_trigger().is_some() || order.is_emulated() {
465            send_emulator_command(TradingCommand::ModifyOrder(command));
466        } else if let Some(algo_id) = order
467            .exec_algorithm_id()
468            .filter(|_| order.is_active_local())
469        {
470            let endpoint = format!("{algo_id}.execute");
471            msgbus::send_any(endpoint.into(), &TradingCommand::ModifyOrder(command));
472        } else {
473            send_risk_command(TradingCommand::ModifyOrder(command));
474        }
475        Ok(())
476    }
477
478    /// Batch modifies multiple orders for the same instrument.
479    ///
480    /// Each tuple is `(client_order_id, quantity, price, trigger_price)`.
481    ///
482    /// # Errors
483    ///
484    /// Returns an error if the strategy is not registered, the orders span multiple instruments,
485    /// contain emulated/local orders, or a child modify is invalid.
486    fn modify_orders(
487        &mut self,
488        updates: Vec<BatchModifyOrder>,
489        client_id: Option<ClientId>,
490        params: Option<Params>,
491    ) -> anyhow::Result<()>
492    where
493        Self: StrategyNative,
494    {
495        if updates.is_empty() {
496            anyhow::bail!("Cannot batch modify empty order list");
497        }
498
499        let (trader_id, strategy_id, ts_init) = {
500            let core = StrategyNative::strategy_core_mut(self);
501            (
502                registered_trader_id(core)?,
503                registered_strategy_id(core)?,
504                core.clock_mut().timestamp_ns(),
505            )
506        };
507
508        let orders: Vec<OrderAny> = {
509            let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
510            let cache = cache_rc.borrow();
511            updates
512                .iter()
513                .map(|(client_order_id, _, _, _)| {
514                    cache
515                        .try_order_owned(client_order_id)
516                        .map_err(|e| anyhow::anyhow!("Cannot modify order: {e}"))
517                })
518                .collect::<Result<_, _>>()?
519        };
520
521        let instrument_id = orders[0].instrument_id();
522
523        for (order, (_, quantity, price, trigger_price)) in orders.iter().zip(updates.iter()) {
524            if order.instrument_id() != instrument_id {
525                anyhow::bail!(
526                    "Cannot batch modify orders for different instruments: {} vs {}",
527                    instrument_id,
528                    order.instrument_id()
529                );
530            }
531
532            if order.is_emulated() || order.is_active_local() {
533                anyhow::bail!("Cannot include emulated or local orders in batch modify");
534            }
535
536            let mut updating = false;
537
538            if quantity.is_some_and(|q| q != order.quantity()) {
539                updating = true;
540            }
541
542            if let Some(price) = price {
543                if !LIMIT_ORDER_TYPES.contains(&order.order_type()) {
544                    anyhow::bail!("{} orders do not have a LIMIT price", order.order_type());
545                }
546
547                if Some(*price) != order.price() {
548                    updating = true;
549                }
550            }
551
552            if let Some(trigger_price) = trigger_price {
553                if !STOP_ORDER_TYPES.contains(&order.order_type()) {
554                    anyhow::bail!(
555                        "{} orders do not have a STOP trigger price",
556                        order.order_type()
557                    );
558                }
559
560                if Some(*trigger_price) != order.trigger_price() {
561                    updating = true;
562                }
563            }
564
565            if !updating {
566                anyhow::bail!(
567                    "Cannot create command BatchModifyOrders: quantity, price, and trigger were \
568                    either None or the same as existing values for {}",
569                    order.client_order_id()
570                );
571            }
572
573            if order.is_closed() || order.is_pending_cancel() {
574                anyhow::bail!(
575                    "Cannot create command BatchModifyOrders: state is {:?}, {order:?}",
576                    order.status()
577                );
578            }
579        }
580
581        let params = params.filter(|params| !params.is_empty());
582        let mut modifies = Vec::with_capacity(orders.len());
583
584        for (order, (_, quantity, price, trigger_price)) in orders.into_iter().zip(updates) {
585            if !self.mark_order_pending_update(&order)? {
586                continue;
587            }
588
589            modifies.push(ModifyOrder::new(
590                trader_id,
591                client_id,
592                strategy_id,
593                instrument_id,
594                order.client_order_id(),
595                order.venue_order_id(),
596                quantity,
597                price,
598                trigger_price,
599                UUID4::new(),
600                ts_init,
601                params.clone(),
602                None, // correlation_id
603            ));
604        }
605
606        if modifies.is_empty() {
607            log::warn!("Cannot send `BatchModifyOrders`, no valid modify commands");
608            return Ok(());
609        }
610
611        let command = BatchModifyOrders::new(
612            trader_id,
613            client_id,
614            strategy_id,
615            instrument_id,
616            modifies,
617            UUID4::new(),
618            ts_init,
619            params,
620            None, // correlation_id
621        );
622
623        send_risk_command(TradingCommand::ModifyOrders(command));
624        Ok(())
625    }
626
627    /// Cancels an order.
628    ///
629    /// # Errors
630    ///
631    /// Returns an error if the strategy is not registered or order cancellation fails.
632    fn cancel_order(
633        &mut self,
634        client_order_id: ClientOrderId,
635        client_id: Option<ClientId>,
636        params: Option<Params>,
637    ) -> anyhow::Result<()>
638    where
639        Self: StrategyNative,
640    {
641        let (trader_id, strategy_id, ts_init) = {
642            let core = StrategyNative::strategy_core_mut(self);
643            (
644                registered_trader_id(core)?,
645                registered_strategy_id(core)?,
646                core.clock_mut().timestamp_ns(),
647            )
648        };
649
650        let params = params.filter(|params| !params.is_empty());
651
652        // TODO: Snapshot the order from the cache. Callers identify it by ID; we own the
653        // snapshot so later calls (which take `&OrderAny` and may re-enter the cache)
654        // run without holding a live cache borrow.
655        let order = StrategyNative::strategy_core_mut(self)
656            .cache_rc()
657            .borrow()
658            .try_order_owned(&client_order_id)
659            .map_err(|e| anyhow::anyhow!("Cannot cancel order: {e}"))?;
660
661        if !self.mark_order_pending_cancel(&order)? {
662            return Ok(());
663        }
664
665        let command = CancelOrder::new(
666            trader_id,
667            client_id,
668            strategy_id,
669            order.instrument_id(),
670            order.client_order_id(),
671            order.venue_order_id(),
672            UUID4::new(),
673            ts_init,
674            params,
675            None, // correlation_id
676        );
677
678        if order.emulation_trigger().is_some() || order.is_emulated() {
679            send_emulator_command(TradingCommand::CancelOrder(command));
680        } else if let Some(algo_id) = order
681            .exec_algorithm_id()
682            .filter(|_| order.is_active_local())
683        {
684            let endpoint = format!("{algo_id}.execute");
685            msgbus::send_any(endpoint.into(), &TradingCommand::CancelOrder(command));
686        } else {
687            send_exec_command(TradingCommand::CancelOrder(command));
688        }
689
690        if StrategyNative::strategy_core(self).config.manage_gtd_expiry
691            && order.time_in_force() == TimeInForce::Gtd
692            && self.has_gtd_expiry_timer(&order.client_order_id())
693        {
694            self.cancel_gtd_expiry(&order.client_order_id());
695        }
696
697        Ok(())
698    }
699
700    /// Batch cancels multiple orders for the same instrument.
701    ///
702    /// # Errors
703    ///
704    /// Returns an error if the strategy is not registered, the orders span multiple instruments,
705    /// or contain emulated/local orders.
706    fn cancel_orders(
707        &mut self,
708        client_order_ids: Vec<ClientOrderId>,
709        client_id: Option<ClientId>,
710        params: Option<Params>,
711    ) -> anyhow::Result<()>
712    where
713        Self: StrategyNative,
714    {
715        if client_order_ids.is_empty() {
716            anyhow::bail!("Cannot batch cancel empty order list");
717        }
718
719        let (trader_id, strategy_id, ts_init) = {
720            let core = StrategyNative::strategy_core_mut(self);
721            (
722                registered_trader_id(core)?,
723                registered_strategy_id(core)?,
724                core.clock_mut().timestamp_ns(),
725            )
726        };
727
728        // TODO: Snapshot all orders from the cache. See `cancel_order` for the rationale.
729        let orders: Vec<OrderAny> = {
730            let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
731            let cache = cache_rc.borrow();
732            client_order_ids
733                .iter()
734                .map(|id| {
735                    cache
736                        .try_order_owned(id)
737                        .map_err(|e| anyhow::anyhow!("Cannot cancel order: {e}"))
738                })
739                .collect::<Result<_, _>>()?
740        };
741
742        let instrument_id = orders[0].instrument_id();
743
744        for order in &orders {
745            if order.instrument_id() != instrument_id {
746                anyhow::bail!(
747                    "Cannot batch cancel orders for different instruments: {} vs {}",
748                    instrument_id,
749                    order.instrument_id()
750                );
751            }
752
753            if order.is_emulated() || order.is_active_local() {
754                anyhow::bail!("Cannot include emulated or local orders in batch cancel");
755            }
756        }
757
758        let mut cancels = Vec::with_capacity(orders.len());
759
760        for order in orders {
761            if !self.mark_order_pending_cancel(&order)? {
762                continue;
763            }
764
765            cancels.push(CancelOrder::new(
766                trader_id,
767                client_id,
768                strategy_id,
769                instrument_id,
770                order.client_order_id(),
771                order.venue_order_id(),
772                UUID4::new(),
773                ts_init,
774                params.clone(),
775                None, // correlation_id
776            ));
777        }
778
779        if cancels.is_empty() {
780            log::warn!("Cannot send `BatchCancelOrders`, no valid cancel commands");
781            return Ok(());
782        }
783
784        let command = BatchCancelOrders::new(
785            trader_id,
786            client_id,
787            strategy_id,
788            instrument_id,
789            cancels,
790            UUID4::new(),
791            ts_init,
792            params,
793            None, // correlation_id
794        );
795
796        send_exec_command(TradingCommand::CancelOrders(command));
797        Ok(())
798    }
799
800    /// Marks an order as pending update locally before the modify command leaves the strategy.
801    ///
802    /// # Errors
803    ///
804    /// Returns an error if applying the pending update event to the cache fails.
805    fn mark_order_pending_update(&mut self, order: &OrderAny) -> anyhow::Result<bool>
806    where
807        Self: StrategyNative,
808    {
809        if order.is_active_local() || order.is_pending_update() {
810            return Ok(true);
811        }
812
813        let strategy_id = order.strategy_id();
814        required_account_id(order, "pending update")?;
815        let event = OrderEventAny::PendingUpdate(self.generate_order_pending_update(order));
816
817        {
818            let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
819            let mut cache = cache_rc.borrow_mut();
820            match cache.update_order(&event) {
821                Ok(_) => {}
822                Err(e)
823                    if matches!(
824                        e.downcast_ref::<OrderError>(),
825                        Some(OrderError::InvalidStateTransition)
826                    ) =>
827                {
828                    log::warn!("InvalidStateTrigger: {e}, did not apply pending update event");
829                    return Ok(false);
830                }
831                Err(e) => return Err(e),
832            }
833        }
834
835        let topic = format!("events.order.{strategy_id}");
836        msgbus::publish_order_event(topic.into(), &event);
837        msgbus::publish_order_event(
838            msgbus::switchboard::get_order_pending_update_topic(order.instrument_id()),
839            &event,
840        );
841
842        Ok(true)
843    }
844
845    /// Marks an order as pending cancel locally before the cancel command leaves the strategy.
846    ///
847    /// # Errors
848    ///
849    /// Returns an error if applying the pending cancel event to the cache fails.
850    fn mark_order_pending_cancel(&mut self, order: &OrderAny) -> anyhow::Result<bool>
851    where
852        Self: StrategyNative,
853    {
854        if order.is_closed() || order.is_pending_cancel() {
855            log::warn!(
856                "Cannot cancel order: state is {:?}, {order:?}",
857                order.status()
858            );
859            return Ok(false);
860        }
861
862        if order.is_active_local() {
863            return Ok(true);
864        }
865
866        let strategy_id = order.strategy_id();
867        required_account_id(order, "pending cancel")?;
868        let event = OrderEventAny::PendingCancel(self.generate_order_pending_cancel(order));
869
870        {
871            let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
872            let mut cache = cache_rc.borrow_mut();
873            match cache.update_order(&event) {
874                Ok(_) => {}
875                Err(e)
876                    if matches!(
877                        e.downcast_ref::<OrderError>(),
878                        Some(OrderError::InvalidStateTransition)
879                    ) =>
880                {
881                    log::warn!("InvalidStateTrigger: {e}, did not apply pending cancel event");
882                    return Ok(false);
883                }
884                Err(e) => return Err(e),
885            }
886            cache.update_order_pending_cancel_local(order);
887        }
888
889        let topic = format!("events.order.{strategy_id}");
890        msgbus::publish_order_event(topic.into(), &event);
891        msgbus::publish_order_event(
892            msgbus::switchboard::get_order_pending_cancel_topic(order.instrument_id()),
893            &event,
894        );
895
896        Ok(true)
897    }
898
899    /// Generates an `OrderPendingUpdate` event for an order.
900    fn generate_order_pending_update(&mut self, order: &OrderAny) -> OrderPendingUpdate
901    where
902        Self: StrategyNative,
903    {
904        let ts_now = StrategyNative::strategy_core_mut(self)
905            .clock_mut()
906            .timestamp_ns();
907
908        OrderPendingUpdate::new(
909            order.trader_id(),
910            order.strategy_id(),
911            order.instrument_id(),
912            order.client_order_id(),
913            order.account_id(),
914            UUID4::new(),
915            ts_now,
916            ts_now,
917            false,
918            order.venue_order_id(),
919        )
920    }
921
922    /// Generates an `OrderPendingCancel` event for an order.
923    fn generate_order_pending_cancel(&mut self, order: &OrderAny) -> OrderPendingCancel
924    where
925        Self: StrategyNative,
926    {
927        let ts_now = StrategyNative::strategy_core_mut(self)
928            .clock_mut()
929            .timestamp_ns();
930
931        OrderPendingCancel::new(
932            order.trader_id(),
933            order.strategy_id(),
934            order.instrument_id(),
935            order.client_order_id(),
936            order.account_id(),
937            UUID4::new(),
938            ts_now,
939            ts_now,
940            false,
941            order.venue_order_id(),
942        )
943    }
944
945    /// Cancels all open orders for the given instrument.
946    ///
947    /// When `strategy_only` is `true`, only orders associated with this strategy are canceled. When
948    /// `false`, one [`CancelAllOrders`] command is sent even when the cache has no matching order.
949    /// The execution engine selects the explicit, venue-routed, or default client and may cancel
950    /// orders associated with other strategies for the same instrument, client, and account. It
951    /// does not broadcast across execution clients.
952    ///
953    /// # Errors
954    ///
955    /// Returns an error if the strategy is not registered or order cancellation fails.
956    fn cancel_all_orders(
957        &mut self,
958        instrument_id: InstrumentId,
959        order_side: Option<OrderSide>,
960        client_id: Option<ClientId>,
961        strategy_only: bool,
962        params: Option<Params>,
963    ) -> anyhow::Result<()>
964    where
965        Self: StrategyNative,
966    {
967        let params = params.filter(|params| !params.is_empty());
968        let core = StrategyNative::strategy_core_mut(self);
969
970        let trader_id = registered_trader_id(core)?;
971        let strategy_id = registered_strategy_id(core)?;
972        let ts_init = core.clock_mut().timestamp_ns();
973
974        if !strategy_only {
975            let command_id = UUID4::new();
976            let command = CancelAllOrders::new(
977                trader_id,
978                client_id,
979                strategy_id,
980                instrument_id,
981                order_side,
982                command_id,
983                ts_init,
984                params,
985                Some(command_id),
986            );
987
988            send_exec_command(TradingCommand::CancelAllOrders(command));
989            return Ok(());
990        }
991
992        let cache = core.cache_ref();
993
994        let mut open_order_ids: Vec<ClientOrderId> = cache
995            .orders_open(
996                None,
997                Some(&instrument_id),
998                Some(&strategy_id),
999                None,
1000                order_side,
1001            )
1002            .into_iter()
1003            .map(|order| order.client_order_id())
1004            .collect();
1005
1006        let mut emulated_order_ids: Vec<ClientOrderId> = cache
1007            .orders_emulated(
1008                None,
1009                Some(&instrument_id),
1010                Some(&strategy_id),
1011                None,
1012                order_side,
1013            )
1014            .into_iter()
1015            .map(|order| order.client_order_id())
1016            .collect();
1017
1018        let mut inflight_order_ids: Vec<ClientOrderId> = cache
1019            .orders_inflight(
1020                None,
1021                Some(&instrument_id),
1022                Some(&strategy_id),
1023                None,
1024                order_side,
1025            )
1026            .into_iter()
1027            .map(|order| order.client_order_id())
1028            .collect();
1029
1030        // Sort the algorithm IDs so the per-algo cancel cascade fires msgbus
1031        // events in a deterministic order across runs; the cache returns an
1032        // unordered AHashSet.
1033        let mut exec_algorithm_ids: Vec<_> = cache.exec_algorithm_ids().into_iter().collect();
1034        exec_algorithm_ids.sort();
1035        let mut algo_order_ids: Vec<ClientOrderId> = Vec::new();
1036
1037        for algo_id in &exec_algorithm_ids {
1038            algo_order_ids.extend(
1039                cache
1040                    .orders_for_exec_algorithm(
1041                        algo_id,
1042                        None,
1043                        Some(&instrument_id),
1044                        Some(&strategy_id),
1045                        None,
1046                        order_side,
1047                    )
1048                    .into_iter()
1049                    .map(|order| order.client_order_id()),
1050            );
1051        }
1052
1053        let matches_client = |client_order_id: &ClientOrderId| {
1054            client_id.is_none_or(|client_id| {
1055                cache
1056                    .client_id(client_order_id)
1057                    .is_none_or(|order_client_id| *order_client_id == client_id)
1058            })
1059        };
1060
1061        open_order_ids.retain(&matches_client);
1062        emulated_order_ids.retain(&matches_client);
1063        inflight_order_ids.retain(&matches_client);
1064        algo_order_ids.retain(&matches_client);
1065
1066        let open_count = open_order_ids.len();
1067        let emulated_count = emulated_order_ids.len();
1068        let inflight_count = inflight_order_ids.len();
1069        let algo_count = algo_order_ids.len();
1070
1071        let mut cancel_routes: Vec<_> = open_order_ids
1072            .iter()
1073            .chain(&emulated_order_ids)
1074            .chain(&inflight_order_ids)
1075            .chain(&algo_order_ids)
1076            .map(|client_order_id| {
1077                (
1078                    *client_order_id,
1079                    client_id.or_else(|| cache.client_id(client_order_id).copied()),
1080                )
1081            })
1082            .collect();
1083        cancel_routes.sort_by_key(|(client_order_id, _)| *client_order_id);
1084        cancel_routes.dedup_by_key(|(client_order_id, _)| *client_order_id);
1085
1086        drop(cache);
1087
1088        if open_count == 0 && emulated_count == 0 && inflight_count == 0 && algo_count == 0 {
1089            let side_str = order_side.map(|s| format!(" {s}")).unwrap_or_default();
1090            log::info!("No {instrument_id} open, emulated, or inflight{side_str} orders to cancel");
1091            return Ok(());
1092        }
1093
1094        let side_str = order_side.map(|s| format!(" {s}")).unwrap_or_default();
1095
1096        if open_count > 0 {
1097            log::info!(
1098                "Canceling {open_count} open{side_str} {instrument_id} order{}",
1099                if open_count == 1 { "" } else { "s" }
1100            );
1101        }
1102
1103        if emulated_count > 0 {
1104            log::info!(
1105                "Canceling {emulated_count} emulated{side_str} {instrument_id} order{}",
1106                if emulated_count == 1 { "" } else { "s" }
1107            );
1108        }
1109
1110        if inflight_count > 0 {
1111            log::info!(
1112                "Canceling {inflight_count} inflight{side_str} {instrument_id} order{}",
1113                if inflight_count == 1 { "" } else { "s" }
1114            );
1115        }
1116
1117        let mut first_error = None;
1118
1119        for (client_order_id, client_id) in cancel_routes {
1120            if let Err(e) = self.cancel_order(client_order_id, client_id, params.clone()) {
1121                if first_error.is_none() {
1122                    first_error = Some(e);
1123                } else {
1124                    log::error!("Error canceling {client_order_id}: {e}");
1125                }
1126            }
1127        }
1128
1129        first_error.map_or(Ok(()), Err)
1130    }
1131
1132    /// Closes a position by submitting a market order for the opposite side.
1133    ///
1134    /// # Errors
1135    ///
1136    /// Returns an error if the strategy is not registered or position closing fails.
1137    #[expect(clippy::too_many_arguments)]
1138    fn close_position(
1139        &mut self,
1140        position: &Position,
1141        client_id: Option<ClientId>,
1142        tags: Option<Vec<Ustr>>,
1143        time_in_force: Option<TimeInForce>,
1144        reduce_only: Option<bool>,
1145        quote_quantity: Option<bool>,
1146        params: Option<Params>,
1147    ) -> anyhow::Result<()>
1148    where
1149        Self: StrategyNative,
1150    {
1151        let core = StrategyNative::strategy_core_mut(self);
1152
1153        if position.is_closed() {
1154            log::warn!("Cannot close position (already closed): {}", position.id);
1155            return Ok(());
1156        }
1157
1158        let Some(closing_side) = OrderCore::closing_side(position.side) else {
1159            log::warn!("Cannot close flat position: {}", position.id);
1160            return Ok(());
1161        };
1162
1163        let order = core.order_factory().market(
1164            position.instrument_id,
1165            closing_side,
1166            position.quantity,
1167            time_in_force,
1168            reduce_only.or(Some(true)),
1169            quote_quantity,
1170            None,
1171            None,
1172            tags,
1173            None,
1174        );
1175
1176        self.submit_order(order, Some(position.id), client_id, params)
1177    }
1178
1179    /// Closes all open positions for the given instrument.
1180    ///
1181    /// # Errors
1182    ///
1183    /// Returns an error if the strategy is not registered or position closing fails.
1184    #[expect(clippy::too_many_arguments)]
1185    fn close_all_positions(
1186        &mut self,
1187        instrument_id: InstrumentId,
1188        position_side: Option<PositionSide>,
1189        client_id: Option<ClientId>,
1190        tags: Option<Vec<Ustr>>,
1191        time_in_force: Option<TimeInForce>,
1192        reduce_only: Option<bool>,
1193        quote_quantity: Option<bool>,
1194        params: Option<Params>,
1195    ) -> anyhow::Result<()>
1196    where
1197        Self: StrategyNative,
1198    {
1199        let core = StrategyNative::strategy_core_mut(self);
1200        let strategy_id = registered_strategy_id(core)?;
1201        let cache = core.cache_ref();
1202
1203        let positions_open = cache.positions_open(
1204            None,
1205            Some(&instrument_id),
1206            Some(&strategy_id),
1207            None,
1208            position_side,
1209        );
1210
1211        let side_str = position_side.map(|s| format!(" {s}")).unwrap_or_default();
1212
1213        if positions_open.is_empty() {
1214            log::info!("No {instrument_id} open{side_str} positions to close");
1215            return Ok(());
1216        }
1217
1218        let count = positions_open.len();
1219        log::info!(
1220            "Closing {count} open{side_str} position{}",
1221            if count == 1 { "" } else { "s" }
1222        );
1223
1224        let positions_data: Vec<_> = positions_open
1225            .iter()
1226            .map(|p| (p.id, p.instrument_id, p.side, p.quantity, p.is_closed()))
1227            .collect();
1228        drop(positions_open);
1229
1230        drop(cache);
1231
1232        for (pos_id, pos_instrument_id, pos_side, pos_quantity, is_closed) in positions_data {
1233            if is_closed {
1234                continue;
1235            }
1236
1237            let core = StrategyNative::strategy_core_mut(self);
1238            let Some(closing_side) = OrderCore::closing_side(pos_side) else {
1239                continue;
1240            };
1241            let order = core.order_factory().market(
1242                pos_instrument_id,
1243                closing_side,
1244                pos_quantity,
1245                time_in_force,
1246                reduce_only.or(Some(true)),
1247                quote_quantity,
1248                None,
1249                None,
1250                tags.clone(),
1251                None,
1252            );
1253
1254            self.submit_order(order, Some(pos_id), client_id, params.clone())?;
1255        }
1256
1257        Ok(())
1258    }
1259
1260    /// Queries account state from the execution client.
1261    ///
1262    /// Creates a [`QueryAccount`] command and sends it to the execution engine,
1263    /// which will request the current account state from the execution client.
1264    ///
1265    /// # Errors
1266    ///
1267    /// Returns an error if the strategy is not registered.
1268    fn query_account(
1269        &mut self,
1270        account_id: AccountId,
1271        client_id: Option<ClientId>,
1272        params: Option<Params>,
1273    ) -> anyhow::Result<()>
1274    where
1275        Self: StrategyNative,
1276    {
1277        let core = StrategyNative::strategy_core_mut(self);
1278
1279        let trader_id = registered_trader_id(core)?;
1280        let ts_init = core.clock_mut().timestamp_ns();
1281
1282        let command = QueryAccount::new(
1283            trader_id,
1284            client_id,
1285            account_id,
1286            UUID4::new(),
1287            ts_init,
1288            params,
1289            None, // correlation_id
1290        );
1291
1292        send_exec_command(TradingCommand::QueryAccount(command));
1293        Ok(())
1294    }
1295
1296    /// Queries order state from the execution client.
1297    ///
1298    /// Creates a [`QueryOrder`] command and sends it to the execution engine,
1299    /// which will request the current order state from the execution client.
1300    ///
1301    /// # Errors
1302    ///
1303    /// Returns an error if the strategy is not registered.
1304    fn query_order(
1305        &mut self,
1306        order: &OrderAny,
1307        client_id: Option<ClientId>,
1308        params: Option<Params>,
1309    ) -> anyhow::Result<()>
1310    where
1311        Self: StrategyNative,
1312    {
1313        let core = StrategyNative::strategy_core_mut(self);
1314
1315        let trader_id = registered_trader_id(core)?;
1316        let strategy_id = registered_strategy_id(core)?;
1317        let ts_init = core.clock_mut().timestamp_ns();
1318
1319        let command = QueryOrder::new(
1320            trader_id,
1321            client_id,
1322            strategy_id,
1323            order.instrument_id(),
1324            order.client_order_id(),
1325            order.venue_order_id(),
1326            UUID4::new(),
1327            ts_init,
1328            params,
1329            None, // correlation_id
1330        );
1331
1332        send_exec_command(TradingCommand::QueryOrder(command));
1333        Ok(())
1334    }
1335
1336    /// Handles an order event, dispatching to the appropriate handler.
1337    fn handle_order_event(&mut self, event: OrderEventAny)
1338    where
1339        Self: StrategyNative,
1340    {
1341        let state = {
1342            let core = StrategyNative::strategy_core_mut(self);
1343            let id = &core.actor.actor_id;
1344            let is_warning = matches!(
1345                &event,
1346                OrderEventAny::Denied(_)
1347                    | OrderEventAny::Rejected(_)
1348                    | OrderEventAny::CancelRejected(_)
1349                    | OrderEventAny::ModifyRejected(_)
1350            );
1351
1352            if is_warning {
1353                log::warn!("{id} {RECV}{EVT} {event}");
1354            } else if core.actor.config.log_events {
1355                log::info!("{id} {RECV}{EVT} {event}");
1356            }
1357
1358            core.actor.state()
1359        };
1360
1361        let client_order_id = event.client_order_id();
1362        let cached_order_is_closed = {
1363            let core = StrategyNative::strategy_core_mut(self);
1364            core.cache_ref()
1365                .order(&client_order_id)
1366                .map(|order| order.is_closed())
1367        };
1368        let is_terminal = match &event {
1369            OrderEventAny::FillVoided(event) => {
1370                cached_order_is_closed.unwrap_or(!event.is_reopened)
1371            }
1372            OrderEventAny::Filled(_) => cached_order_is_closed.unwrap_or(true),
1373            OrderEventAny::Canceled(_)
1374            | OrderEventAny::Rejected(_)
1375            | OrderEventAny::Expired(_)
1376            | OrderEventAny::Denied(_) => true,
1377            _ => false,
1378        };
1379
1380        // GTD timer cleanup runs regardless of state so timers do not leak when
1381        // terminal events arrive during the post-stop delay.
1382        if is_terminal {
1383            self.cancel_gtd_expiry(&client_order_id);
1384        }
1385
1386        // Events are logged unconditionally so residual events received after stop
1387        // remain observable, but dispatch is gated on the running state.
1388        if state != ComponentState::Running {
1389            return;
1390        }
1391
1392        if matches!(&event, OrderEventAny::FillVoided(event) if event.is_reopened) {
1393            let order = StrategyNative::strategy_core_mut(self)
1394                .cache_ref()
1395                .order(&client_order_id)
1396                .map(|order| order.clone());
1397            if let Some(order) = order
1398                && order.is_open()
1399                && !self.has_gtd_expiry_timer(&client_order_id)
1400                && let Err(e) = self.set_gtd_expiry(&order)
1401            {
1402                log::error!(
1403                    "Failed to restore GTD expiry for reopened order {client_order_id}: {e}"
1404                );
1405            }
1406        }
1407
1408        let manager_actions = {
1409            let core = StrategyNative::strategy_core_mut(self);
1410            if core.config.manage_contingent_orders {
1411                core.order_manager
1412                    .as_mut()
1413                    .map_or_else(Vec::new, |manager| manager.handle_event(&event))
1414            } else {
1415                Vec::new()
1416            }
1417        };
1418        self.dispatch_manager_actions(manager_actions);
1419
1420        match &event {
1421            OrderEventAny::Initialized(e) => self.on_order_initialized(e.clone()),
1422            OrderEventAny::Denied(e) => self.on_order_denied(*e),
1423            OrderEventAny::Emulated(e) => self.on_order_emulated(*e),
1424            OrderEventAny::Released(e) => self.on_order_released(*e),
1425            OrderEventAny::Submitted(e) => self.on_order_submitted(*e),
1426            OrderEventAny::Rejected(e) => self.on_order_rejected(*e),
1427            OrderEventAny::Accepted(e) => self.on_order_accepted(*e),
1428            OrderEventAny::Canceled(e) => self.on_order_canceled(e),
1429            OrderEventAny::Expired(e) => self.on_order_expired(*e),
1430            OrderEventAny::Triggered(e) => self.on_order_triggered(*e),
1431            OrderEventAny::PendingUpdate(e) => self.on_order_pending_update(*e),
1432            OrderEventAny::PendingCancel(e) => self.on_order_pending_cancel(*e),
1433            OrderEventAny::ModifyRejected(e) => self.on_order_modify_rejected(*e),
1434            OrderEventAny::CancelRejected(e) => self.on_order_cancel_rejected(*e),
1435            OrderEventAny::Updated(e) => self.on_order_updated(*e),
1436            OrderEventAny::Filled(e) => self.on_order_filled(e),
1437            OrderEventAny::FillVoided(e) => self.on_order_fill_voided(e),
1438        }
1439        self.on_order_event(event);
1440    }
1441
1442    fn dispatch_manager_actions(&mut self, actions: Vec<OrderManagerAction>)
1443    where
1444        Self: StrategyNative,
1445    {
1446        for action in actions {
1447            match action {
1448                OrderManagerAction::PublishInitialized(event) => {
1449                    let topic = msgbus::switchboard::get_event_order_topic(event.strategy_id());
1450                    msgbus::publish_order_event(topic, &event);
1451                }
1452                OrderManagerAction::SubmitToEmulator(command) => {
1453                    send_emulator_command(TradingCommand::SubmitOrder(command));
1454                }
1455                OrderManagerAction::SubmitToRisk(command) => {
1456                    send_risk_command(TradingCommand::SubmitOrder(command));
1457                }
1458                OrderManagerAction::SubmitToAlgorithm {
1459                    command,
1460                    exec_algorithm_id,
1461                } => send_algo_command(command, exec_algorithm_id),
1462                OrderManagerAction::CancelLocal(order) => {
1463                    let client_order_id = order.client_order_id();
1464                    if let Err(e) = self.cancel_order(client_order_id, None, None) {
1465                        log::error!(
1466                            "Failed to dispatch contingent cancel for {client_order_id}: {e}"
1467                        );
1468                    }
1469                }
1470                OrderManagerAction::ModifyLocalQuantity { order, quantity } => {
1471                    let client_order_id = order.client_order_id();
1472                    if let Err(e) =
1473                        self.modify_order(client_order_id, Some(quantity), None, None, None, None)
1474                    {
1475                        log::error!(
1476                            "Failed to dispatch contingent modify for {client_order_id}: {e}"
1477                        );
1478                    }
1479                }
1480            }
1481        }
1482    }
1483
1484    /// Handles a position event, dispatching to the appropriate handler.
1485    fn handle_position_event(&mut self, event: PositionEvent)
1486    where
1487        Self: StrategyNative,
1488    {
1489        let state = {
1490            let core = StrategyNative::strategy_core_mut(self);
1491
1492            if core.actor.config.log_events {
1493                let id = &core.actor.actor_id;
1494                log::info!("{id} {RECV}{EVT} {event:?}");
1495            }
1496
1497            core.actor.state()
1498        };
1499
1500        if state != ComponentState::Running {
1501            return;
1502        }
1503
1504        match &event {
1505            PositionEvent::PositionOpened(e) => self.on_position_opened(e.clone()),
1506            PositionEvent::PositionChanged(e) => self.on_position_changed(e.clone()),
1507            PositionEvent::PositionClosed(e) => self.on_position_closed(e.clone()),
1508            PositionEvent::PositionAdjusted(_) => {
1509                return;
1510            }
1511        }
1512        self.on_position_event(event);
1513    }
1514
1515    // -- LIFECYCLE METHODS -----------------------------------------------------------------------
1516
1517    /// Called when the strategy is started.
1518    ///
1519    /// Override this method to implement custom initialization logic.
1520    /// The default implementation reactivates GTD timers if `manage_gtd_expiry` is enabled.
1521    ///
1522    /// # Errors
1523    ///
1524    /// Returns an error if strategy initialization fails.
1525    fn on_start(&mut self) -> anyhow::Result<()>
1526    where
1527        Self: StrategyNative,
1528    {
1529        let core = StrategyNative::strategy_core_mut(self);
1530        let strategy_id = registered_strategy_id(core)?;
1531        log::info!("Starting {strategy_id}");
1532
1533        if core.config.manage_gtd_expiry {
1534            self.reactivate_gtd_timers();
1535        }
1536
1537        Ok(())
1538    }
1539
1540    /// Routes a time event through the framework-managed strategy handlers when called directly.
1541    ///
1542    /// Framework timer dispatch never calls this method. Implement [`DataActor::on_time_event`] for
1543    /// user timer callbacks. Runtime hosts call [`route_time_event`] directly so an override cannot
1544    /// replace framework-managed GTD expiry or market exit routing.
1545    ///
1546    /// # Errors
1547    ///
1548    /// Returns an error if time event handling fails.
1549    fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()>
1550    where
1551        Self: StrategyNative + Component,
1552    {
1553        route_time_event(self, event);
1554        Ok(())
1555    }
1556
1557    // -- EVENT HANDLERS --------------------------------------------------------------------------
1558
1559    /// Called when an order is initialized.
1560    ///
1561    /// Override this method to implement custom logic when an order is first created.
1562    #[allow(unused_variables)]
1563    fn on_order_initialized(&mut self, event: OrderInitialized) {}
1564
1565    /// Called when any order event is received after the specific order handler runs.
1566    ///
1567    /// Override this method to implement custom logic for all order events.
1568    #[allow(unused_variables)]
1569    fn on_order_event(&mut self, event: OrderEventAny) {}
1570
1571    /// Called when an order is denied by the system.
1572    ///
1573    /// Override this method to implement custom logic when an order is denied before submission.
1574    #[allow(unused_variables)]
1575    fn on_order_denied(&mut self, event: OrderDenied) {}
1576
1577    /// Called when an order is emulated.
1578    ///
1579    /// Override this method to implement custom logic when an order is taken over by the emulator.
1580    #[allow(unused_variables)]
1581    fn on_order_emulated(&mut self, event: OrderEmulated) {}
1582
1583    /// Called when an order is released from emulation.
1584    ///
1585    /// Override this method to implement custom logic when an emulated order is released.
1586    #[allow(unused_variables)]
1587    fn on_order_released(&mut self, event: OrderReleased) {}
1588
1589    /// Called when an order is submitted to the venue.
1590    ///
1591    /// Override this method to implement custom logic when an order is submitted.
1592    #[allow(unused_variables)]
1593    fn on_order_submitted(&mut self, event: OrderSubmitted) {}
1594
1595    /// Called when an order is rejected by the venue.
1596    ///
1597    /// Override this method to implement custom logic when an order is rejected.
1598    #[allow(unused_variables)]
1599    fn on_order_rejected(&mut self, event: OrderRejected) {}
1600
1601    /// Called when an order is accepted by the venue.
1602    ///
1603    /// Override this method to implement custom logic when an order is accepted.
1604    #[allow(unused_variables)]
1605    fn on_order_accepted(&mut self, event: OrderAccepted) {}
1606
1607    /// Called when an order expires.
1608    ///
1609    /// Override this method to implement custom logic when an order expires.
1610    #[allow(unused_variables)]
1611    fn on_order_expired(&mut self, event: OrderExpired) {}
1612
1613    /// Called when an order is triggered.
1614    ///
1615    /// Override this method to implement custom logic when a stop or conditional order is triggered.
1616    #[allow(unused_variables)]
1617    fn on_order_triggered(&mut self, event: OrderTriggered) {}
1618
1619    /// Called when an order modification is pending.
1620    ///
1621    /// Override this method to implement custom logic when an order is pending modification.
1622    #[allow(unused_variables)]
1623    fn on_order_pending_update(&mut self, event: OrderPendingUpdate) {}
1624
1625    /// Called when an order cancellation is pending.
1626    ///
1627    /// Override this method to implement custom logic when an order is pending cancellation.
1628    #[allow(unused_variables)]
1629    fn on_order_pending_cancel(&mut self, event: OrderPendingCancel) {}
1630
1631    /// Called when an order modification is rejected.
1632    ///
1633    /// Override this method to implement custom logic when an order modification is rejected.
1634    #[allow(unused_variables)]
1635    fn on_order_modify_rejected(&mut self, event: OrderModifyRejected) {}
1636
1637    /// Called when an order cancellation is rejected.
1638    ///
1639    /// Override this method to implement custom logic when an order cancellation is rejected.
1640    #[allow(unused_variables)]
1641    fn on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {}
1642
1643    /// Called when an order is updated.
1644    ///
1645    /// Override this method to implement custom logic when an order is modified.
1646    #[allow(unused_variables)]
1647    fn on_order_updated(&mut self, event: OrderUpdated) {}
1648
1649    /// Called when an order is canceled.
1650    ///
1651    /// Override this method to implement custom logic when an order is canceled.
1652    #[allow(unused_variables)]
1653    fn on_order_canceled(&mut self, event: &OrderCanceled) {}
1654
1655    /// Called when an order is filled.
1656    ///
1657    /// Override this method to implement custom logic when an order is filled.
1658    #[allow(unused_variables)]
1659    fn on_order_filled(&mut self, event: &OrderFilled) {}
1660
1661    /// Called when an applied order fill is partly or fully voided.
1662    #[allow(unused_variables)]
1663    fn on_order_fill_voided(&mut self, event: &OrderFillVoided) {}
1664
1665    /// Called when a position is opened.
1666    ///
1667    /// Override this method to implement custom logic when a position is opened.
1668    #[allow(unused_variables)]
1669    fn on_position_opened(&mut self, event: PositionOpened) {}
1670
1671    /// Called after a position opened, changed, or closed handler runs.
1672    ///
1673    /// Override this method to implement custom logic for all position events.
1674    #[allow(unused_variables)]
1675    fn on_position_event(&mut self, event: PositionEvent) {}
1676
1677    /// Called when a position is changed (quantity or price updated).
1678    ///
1679    /// Override this method to implement custom logic when a position changes.
1680    #[allow(unused_variables)]
1681    fn on_position_changed(&mut self, event: PositionChanged) {}
1682
1683    /// Called when a position is closed.
1684    ///
1685    /// Override this method to implement custom logic when a position is closed.
1686    #[allow(unused_variables)]
1687    fn on_position_closed(&mut self, event: PositionClosed) {}
1688
1689    /// Called when a market exit has been initiated.
1690    ///
1691    /// Override this method to implement custom logic when a market exit begins.
1692    fn on_market_exit(&mut self) {}
1693
1694    /// Called after a market exit has completed.
1695    ///
1696    /// Override this method to implement custom logic after a market exit completes.
1697    fn post_market_exit(&mut self) {}
1698
1699    /// Returns whether the strategy is currently executing a market exit.
1700    ///
1701    /// Strategies can check this to avoid submitting new orders during exit.
1702    fn is_exiting(&self) -> bool
1703    where
1704        Self: StrategyNative,
1705    {
1706        StrategyNative::strategy_core(self).is_exiting
1707    }
1708
1709    /// Initiates an iterative market exit for the strategy.
1710    ///
1711    /// Will cancel all open orders and close all open positions, and wait for
1712    /// all in-flight orders to resolve and positions to close. The strategy
1713    /// remains running after the exit completes.
1714    ///
1715    /// The `on_market_exit` hook is called when the exit process begins.
1716    /// The `post_market_exit` hook is called when the exit process completes.
1717    ///
1718    /// Uses `market_exit_time_in_force` and `market_exit_reduce_only` from
1719    /// the strategy config for closing market orders.
1720    ///
1721    /// # Errors
1722    ///
1723    /// Returns an error if the market exit cannot be initiated.
1724    fn market_exit(&mut self) -> anyhow::Result<()>
1725    where
1726        Self: StrategyNative,
1727    {
1728        let core = StrategyNative::strategy_core_mut(self);
1729        let strategy_id = registered_strategy_id(core)?;
1730
1731        if core.actor.state() != ComponentState::Running {
1732            log::warn!("{strategy_id} Cannot market exit: strategy is not running");
1733            return Ok(());
1734        }
1735
1736        if core.is_exiting {
1737            log::warn!("{strategy_id} Market exit called when already in progress");
1738            return Ok(());
1739        }
1740
1741        core.is_exiting = true;
1742        core.market_exit_attempts = 0;
1743        let time_in_force = core.config.market_exit_time_in_force;
1744        let reduce_only = core.config.market_exit_reduce_only;
1745
1746        log::info!("{strategy_id} Initiating market exit...");
1747
1748        self.on_market_exit();
1749
1750        let core = StrategyNative::strategy_core_mut(self);
1751        let cache = core.cache_ref();
1752
1753        let mut instruments: AHashSet<InstrumentId> = AHashSet::new();
1754
1755        for client_order_id in
1756            cache.iter_client_order_ids_open(None, None, Some(&strategy_id), None)
1757        {
1758            if let Some(order) = cache.order(&client_order_id) {
1759                instruments.insert(order.instrument_id());
1760            }
1761        }
1762
1763        for client_order_id in
1764            cache.iter_client_order_ids_inflight(None, None, Some(&strategy_id), None)
1765        {
1766            if let Some(order) = cache.order(&client_order_id) {
1767                instruments.insert(order.instrument_id());
1768            }
1769        }
1770
1771        for position_id in cache.iter_position_open_ids(None, None, Some(&strategy_id), None) {
1772            if let Some(position) = cache.position(&position_id) {
1773                instruments.insert(position.instrument_id);
1774            }
1775        }
1776
1777        let market_exit_tag = core.market_exit_tag;
1778        // Sort so the per-instrument cancel_all_orders/close_all_positions
1779        // cascade fires msgbus commands in a deterministic sequence; the
1780        // upstream dedup is AHash-backed.
1781        let mut instruments: Vec<_> = instruments.into_iter().collect();
1782        instruments.sort();
1783        drop(cache);
1784
1785        for instrument_id in instruments {
1786            if let Err(e) = self.cancel_all_orders(instrument_id, None, None, true, None) {
1787                log::error!("Error canceling orders for {instrument_id}: {e}");
1788            }
1789
1790            if let Err(e) = self.close_all_positions(
1791                instrument_id,
1792                None,
1793                None,
1794                Some(vec![market_exit_tag]),
1795                Some(time_in_force),
1796                Some(reduce_only),
1797                None,
1798                None,
1799            ) {
1800                log::error!("Error closing positions for {instrument_id}: {e}");
1801            }
1802        }
1803
1804        let core = StrategyNative::strategy_core_mut(self);
1805        let interval_ms = core.config.market_exit_interval_ms;
1806        let timer_name = core.market_exit_timer_name;
1807
1808        log::info!("{strategy_id} Setting market exit timer at {interval_ms}ms intervals");
1809
1810        let interval_ns = interval_ms * 1_000_000;
1811        let result = core.clock_mut().set_timer_ns(
1812            timer_name.as_str(),
1813            interval_ns,
1814            None,
1815            None,
1816            None,
1817            None,
1818            None,
1819        );
1820
1821        if let Err(e) = result {
1822            // Reset exit state on timer failure (caller handles pending_stop)
1823            core.is_exiting = false;
1824            core.market_exit_attempts = 0;
1825            return Err(e);
1826        }
1827
1828        Ok(())
1829    }
1830
1831    /// Checks if the market exit is complete and finalizes if so.
1832    ///
1833    /// This method is called by the market exit timer.
1834    fn check_market_exit(&mut self, _event: TimeEvent)
1835    where
1836        Self: StrategyNative + Component,
1837    {
1838        // Guard against stale timer events after cancel_market_exit
1839        if !self.is_exiting() {
1840            return;
1841        }
1842
1843        let core = StrategyNative::strategy_core_mut(self);
1844        let Some(strategy_id) = core.strategy_id() else {
1845            log::error!("Cannot check market exit: strategy_id is not set");
1846            return;
1847        };
1848
1849        core.market_exit_attempts += 1;
1850        let attempts = core.market_exit_attempts;
1851        let max_attempts = core.config.market_exit_max_attempts;
1852
1853        log::debug!(
1854            "{strategy_id} Market exit check triggered (attempt {attempts}/{max_attempts})"
1855        );
1856
1857        if attempts >= max_attempts {
1858            let cache = core.cache_ref();
1859            let open_orders_count =
1860                cache.orders_open_count(None, None, Some(&strategy_id), None, None);
1861            let inflight_orders_count =
1862                cache.orders_inflight_count(None, None, Some(&strategy_id), None, None);
1863            let open_positions_count =
1864                cache.positions_open_count(None, None, Some(&strategy_id), None, None);
1865
1866            drop(cache);
1867
1868            log::warn!(
1869                "{strategy_id} Market exit max attempts ({max_attempts}) reached, \
1870                completing with open orders: {open_orders_count}, \
1871                inflight orders: {inflight_orders_count}, \
1872                open positions: {open_positions_count}"
1873            );
1874
1875            self.finalize_market_exit();
1876            return;
1877        }
1878
1879        let cache = core.cache_ref();
1880        let has_open_orders = !cache
1881            .orders_open(None, None, Some(&strategy_id), None, None)
1882            .is_empty();
1883        let has_inflight_orders = !cache
1884            .orders_inflight(None, None, Some(&strategy_id), None, None)
1885            .is_empty();
1886
1887        if has_open_orders || has_inflight_orders {
1888            return;
1889        }
1890
1891        let positions_data: Vec<_> = cache
1892            .positions_open(None, None, Some(&strategy_id), None, None)
1893            .iter()
1894            .map(|p| (p.id, p.instrument_id, p.side, p.quantity, p.is_closed()))
1895            .collect();
1896
1897        if !positions_data.is_empty() {
1898            // If there are open positions but no orders, re-send close orders
1899            drop(cache);
1900
1901            for (pos_id, instrument_id, side, quantity, is_closed) in positions_data {
1902                if is_closed {
1903                    continue;
1904                }
1905
1906                let core = StrategyNative::strategy_core_mut(self);
1907                let time_in_force = core.config.market_exit_time_in_force;
1908                let reduce_only = core.config.market_exit_reduce_only;
1909                let market_exit_tag = core.market_exit_tag;
1910                let Some(closing_side) = OrderCore::closing_side(side) else {
1911                    continue;
1912                };
1913                let order = core.order_factory().market(
1914                    instrument_id,
1915                    closing_side,
1916                    quantity,
1917                    Some(time_in_force),
1918                    Some(reduce_only),
1919                    None,
1920                    None,
1921                    None,
1922                    Some(vec![market_exit_tag]),
1923                    None,
1924                );
1925
1926                if let Err(e) = self.submit_order(order, Some(pos_id), None, None) {
1927                    log::error!("Error re-submitting close order for position {pos_id}: {e}");
1928                }
1929            }
1930            return;
1931        }
1932
1933        drop(cache);
1934        self.finalize_market_exit();
1935    }
1936
1937    /// Finalizes the market exit process.
1938    ///
1939    /// Cancels the market exit timer, resets state, calls the `post_market_exit` hook,
1940    /// and stops the strategy if a stop was pending.
1941    fn finalize_market_exit(&mut self)
1942    where
1943        Self: StrategyNative + Component,
1944    {
1945        let (actor_id, should_stop) = {
1946            let core = StrategyNative::strategy_core_mut(self);
1947            let actor_id = core.actor_id();
1948            let should_stop = core.pending_stop;
1949            (actor_id, should_stop)
1950        };
1951
1952        self.cancel_market_exit();
1953
1954        let hook_result = catch_unwind(AssertUnwindSafe(|| {
1955            self.post_market_exit();
1956        }));
1957
1958        if let Err(e) = hook_result {
1959            log::error!("{actor_id} Error in post_market_exit: {e:?}");
1960        }
1961
1962        if should_stop {
1963            log::info!("{actor_id} Market exit complete, stopping strategy");
1964
1965            if let Err(e) = Component::stop(self) {
1966                log::error!("{actor_id} Failed to stop: {e}");
1967            }
1968        }
1969
1970        let core = StrategyNative::strategy_core_mut(self);
1971        debug_assert!(
1972            !(core.pending_stop
1973                && !core.is_exiting
1974                && core.actor.state() == ComponentState::Running),
1975            "INVARIANT: stuck state after finalize_market_exit"
1976        );
1977    }
1978
1979    /// Cancels an active market exit without calling hooks.
1980    ///
1981    /// Used when `stop()` is called during an active market exit to avoid state leaks.
1982    fn cancel_market_exit(&mut self)
1983    where
1984        Self: StrategyNative,
1985    {
1986        let core = StrategyNative::strategy_core_mut(self);
1987        let timer_name = core.market_exit_timer_name;
1988
1989        if core
1990            .clock_mut()
1991            .timer_names()
1992            .contains(&timer_name.as_str())
1993        {
1994            core.clock_mut().cancel_timer(timer_name.as_str());
1995        }
1996
1997        core.is_exiting = false;
1998        core.pending_stop = false;
1999        core.market_exit_attempts = 0;
2000    }
2001
2002    /// Stops the strategy with optional managed stop behavior.
2003    ///
2004    /// If `manage_stop` is enabled in the config, the strategy will first complete
2005    /// any active market exit (or initiate one) before stopping. If `manage_stop`
2006    /// is disabled, the strategy stops immediately, cleaning up any active market
2007    /// exit state.
2008    ///
2009    /// # Returns
2010    ///
2011    /// Returns `true` if the strategy should proceed with stopping, `false` if
2012    /// the stop is being deferred until market exit completes.
2013    fn stop(&mut self) -> bool
2014    where
2015        Self: StrategyNative,
2016    {
2017        let (manage_stop, is_exiting, should_initiate_exit) = {
2018            let core = StrategyNative::strategy_core_mut(self);
2019            let actor_id = core.actor_id();
2020            let manage_stop = core.config.manage_stop;
2021            let state = core.actor.state();
2022            let pending_stop = core.pending_stop;
2023            let is_exiting = core.is_exiting;
2024
2025            if manage_stop {
2026                if state != ComponentState::Running {
2027                    return true; // Proceed with stop
2028                }
2029
2030                if pending_stop {
2031                    return false; // Already waiting for market exit
2032                }
2033
2034                core.pending_stop = true;
2035                let should_initiate_exit = !is_exiting;
2036
2037                if should_initiate_exit {
2038                    log::info!("{actor_id} Initiating market exit before stop");
2039                }
2040
2041                (manage_stop, is_exiting, should_initiate_exit)
2042            } else {
2043                (manage_stop, is_exiting, false)
2044            }
2045        };
2046
2047        if manage_stop {
2048            if should_initiate_exit && let Err(e) = self.market_exit() {
2049                log::warn!("Market exit failed during stop: {e}, proceeding with stop");
2050                StrategyNative::strategy_core_mut(self).pending_stop = false;
2051                return true;
2052            }
2053            debug_assert!(
2054                self.is_exiting(),
2055                "INVARIANT: deferring stop but not exiting"
2056            );
2057            return false; // Defer stop until market exit completes
2058        }
2059
2060        // manage_stop is false - clean up any active market exit
2061        if is_exiting {
2062            self.cancel_market_exit();
2063        }
2064
2065        true // Proceed with stop
2066    }
2067
2068    /// Denies an order by generating an `OrderDenied` event.
2069    ///
2070    /// This method creates an `OrderDenied` event, applies it to the order,
2071    /// and updates the cache.
2072    fn deny_order(&mut self, order: &OrderAny, reason: Ustr)
2073    where
2074        Self: StrategyNative,
2075    {
2076        let core = StrategyNative::strategy_core_mut(self);
2077        let Some(trader_id) = core.trader_id() else {
2078            log::error!(
2079                "Cannot deny order {}: trader_id is not set",
2080                order.client_order_id()
2081            );
2082            return;
2083        };
2084        let Some(strategy_id) = core.strategy_id() else {
2085            log::error!(
2086                "Cannot deny order {}: strategy_id is not set",
2087                order.client_order_id()
2088            );
2089            return;
2090        };
2091        let ts_now = core.clock_mut().timestamp_ns();
2092
2093        let event = OrderDenied::new(
2094            trader_id,
2095            strategy_id,
2096            order.instrument_id(),
2097            order.client_order_id(),
2098            reason,
2099            UUID4::new(),
2100            ts_now,
2101            ts_now,
2102        );
2103
2104        log::warn!(
2105            "{strategy_id} Order {} denied: {reason}",
2106            order.client_order_id()
2107        );
2108
2109        let publish_initialized = {
2110            let cache_rc = core.cache_rc();
2111            let mut cache = cache_rc.borrow_mut();
2112            if cache.order_exists(&order.client_order_id()) {
2113                false
2114            } else {
2115                match cache.add_order(order.clone(), None, None, true) {
2116                    Ok(()) => true,
2117                    Err(e) => {
2118                        log::warn!("Failed to add denied order to cache: {e}");
2119                        false
2120                    }
2121                }
2122            }
2123        };
2124
2125        if publish_initialized {
2126            publish_order_initialized(order);
2127        }
2128
2129        let event = OrderEventAny::Denied(event);
2130        let applied = {
2131            let cache_rc = core.cache_rc();
2132            let mut cache = cache_rc.borrow_mut();
2133            if let Err(e) = cache.update_order(&event) {
2134                log::warn!("Failed to apply OrderDenied event: {e}");
2135                false
2136            } else {
2137                true
2138            }
2139        };
2140
2141        if applied {
2142            let topic = format!("events.order.{strategy_id}");
2143            msgbus::publish_order_event(topic.into(), &event);
2144        }
2145    }
2146
2147    /// Denies all orders in an order list.
2148    ///
2149    /// This method denies each non-closed order in the list.
2150    fn deny_order_list(&mut self, orders: &[OrderAny], reason: Ustr)
2151    where
2152        Self: StrategyNative,
2153    {
2154        for order in orders {
2155            if !order.is_closed() {
2156                self.deny_order(order, reason);
2157            }
2158        }
2159    }
2160
2161    // -- GTD EXPIRY MANAGEMENT -------------------------------------------------------------------
2162
2163    /// Sets a GTD expiry timer for an order.
2164    ///
2165    /// Creates a timer that will automatically cancel the order when it expires.
2166    ///
2167    /// # Errors
2168    ///
2169    /// Returns an error if timer creation fails.
2170    fn set_gtd_expiry(&mut self, order: &OrderAny) -> anyhow::Result<()>
2171    where
2172        Self: StrategyNative,
2173    {
2174        let core = StrategyNative::strategy_core_mut(self);
2175
2176        if !core.config.manage_gtd_expiry || order.time_in_force() != TimeInForce::Gtd {
2177            return Ok(());
2178        }
2179
2180        let Some(expire_time) = order.expire_time() else {
2181            return Ok(());
2182        };
2183
2184        let client_order_id = order.client_order_id();
2185        let timer_name = format!("GTD-EXPIRY:{client_order_id}");
2186
2187        let current_time_ns = {
2188            let clock = core.clock_mut();
2189            clock.timestamp_ns()
2190        };
2191
2192        if current_time_ns >= expire_time.as_u64() {
2193            log::info!("GTD order {client_order_id} already expired, canceling immediately");
2194            return self.cancel_order(order.client_order_id(), None, None);
2195        }
2196
2197        {
2198            let mut clock = core.clock_mut();
2199            clock.set_time_alert_ns(&timer_name, expire_time, None, None)?;
2200        }
2201
2202        core.gtd_timers
2203            .insert(client_order_id, Ustr::from(&timer_name));
2204
2205        log::debug!("Set GTD expiry timer for {client_order_id} at {expire_time}");
2206        Ok(())
2207    }
2208
2209    /// Cancels a GTD expiry timer for an order.
2210    fn cancel_gtd_expiry(&mut self, client_order_id: &ClientOrderId)
2211    where
2212        Self: StrategyNative,
2213    {
2214        let core = StrategyNative::strategy_core_mut(self);
2215
2216        if let Some(timer_name) = core.gtd_timers.remove(client_order_id) {
2217            core.clock_mut().cancel_timer(timer_name.as_str());
2218            log::debug!("Canceled GTD expiry timer for {client_order_id}");
2219        }
2220    }
2221
2222    /// Checks if a GTD expiry timer exists for an order.
2223    fn has_gtd_expiry_timer(&mut self, client_order_id: &ClientOrderId) -> bool
2224    where
2225        Self: StrategyNative,
2226    {
2227        let core = StrategyNative::strategy_core_mut(self);
2228        core.gtd_timers.contains_key(client_order_id)
2229    }
2230
2231    /// Handles GTD order expiry by canceling the order.
2232    ///
2233    /// This method is called when a GTD expiry timer fires.
2234    fn expire_gtd_order(&mut self, event: TimeEvent)
2235    where
2236        Self: StrategyNative,
2237    {
2238        let timer_name = event.name;
2239        let Some(client_order_id) = timer_name
2240            .as_str()
2241            .strip_prefix("GTD-EXPIRY:")
2242            .and_then(|value| ClientOrderId::new_checked(value).ok())
2243        else {
2244            log::error!("Invalid GTD timer name format: {timer_name}");
2245            return;
2246        };
2247
2248        let core = StrategyNative::strategy_core_mut(self);
2249        if core.gtd_timers.get(&client_order_id) != Some(&timer_name) {
2250            return;
2251        }
2252        core.gtd_timers.remove(&client_order_id);
2253
2254        let order = core.cache_ref().order(&client_order_id).map(|o| o.clone());
2255        let Some(order) = order else {
2256            log::warn!("GTD order {client_order_id} not found in cache");
2257            return;
2258        };
2259
2260        log::info!("GTD order {client_order_id} expired");
2261
2262        if let Err(e) = self.cancel_order(order.client_order_id(), None, None) {
2263            log::error!("Failed to cancel expired GTD order {client_order_id}: {e}");
2264        }
2265    }
2266
2267    /// Reactivates GTD timers for open orders on strategy start.
2268    ///
2269    /// Queries the cache for all open GTD orders and creates timers for those
2270    /// that haven't expired yet. Orders that have already expired are canceled immediately.
2271    fn reactivate_gtd_timers(&mut self)
2272    where
2273        Self: StrategyNative,
2274    {
2275        let core = StrategyNative::strategy_core_mut(self);
2276        let Some(strategy_id) = core.strategy_id() else {
2277            log::error!("Cannot reactivate GTD timers: strategy_id is not set");
2278            return;
2279        };
2280        let current_time_ns = core.clock_mut().timestamp_ns();
2281
2282        let gtd_orders: Vec<OrderAny> = core
2283            .cache_ref()
2284            .orders_open(None, None, Some(&strategy_id), None, None)
2285            .into_iter()
2286            .filter(|o| o.time_in_force() == TimeInForce::Gtd)
2287            .map(|o| o.clone())
2288            .collect();
2289
2290        for order in gtd_orders {
2291            let Some(expire_time) = order.expire_time() else {
2292                continue;
2293            };
2294
2295            let expire_time_ns = expire_time.as_u64();
2296            let client_order_id = order.client_order_id();
2297
2298            if current_time_ns >= expire_time_ns {
2299                log::info!("GTD order {client_order_id} already expired, canceling immediately");
2300                if let Err(e) = self.cancel_order(order.client_order_id(), None, None) {
2301                    log::error!("Failed to cancel expired GTD order {client_order_id}: {e}");
2302                }
2303            } else if let Err(e) = self.set_gtd_expiry(&order) {
2304                log::error!("Failed to set GTD expiry timer for {client_order_id}: {e}");
2305            }
2306        }
2307    }
2308}
2309
2310/// Routes a time event through framework-managed strategy handlers.
2311///
2312/// Runtime hosts call this before the user-facing [`DataActor::on_time_event`] callback so custom
2313/// [`Strategy::on_time_event`] implementations cannot replace managed GTD expiry or market exit
2314/// routing.
2315pub fn route_time_event<T>(strategy: &mut T, event: &TimeEvent)
2316where
2317    T: Strategy + StrategyNative + Component + ?Sized,
2318{
2319    let (gtd_order_id, is_market_exit) = {
2320        let core = StrategyNative::strategy_core(strategy);
2321        let gtd_order_id = event
2322            .name
2323            .as_str()
2324            .strip_prefix("GTD-EXPIRY:")
2325            .and_then(|value| ClientOrderId::new_checked(value).ok())
2326            .filter(|client_order_id| core.gtd_timers.get(client_order_id) == Some(&event.name));
2327        let is_market_exit = event.name == core.market_exit_timer_name;
2328        (gtd_order_id, is_market_exit)
2329    };
2330
2331    if gtd_order_id.is_none() && !is_market_exit {
2332        return;
2333    }
2334
2335    let core = StrategyNative::strategy_core_mut(strategy);
2336    if core.managed_time_event_last_id == Some(event.event_id) {
2337        return;
2338    }
2339    core.managed_time_event_last_id = Some(event.event_id);
2340
2341    if gtd_order_id.is_some() {
2342        strategy.expire_gtd_order(event.clone());
2343    } else {
2344        strategy.check_market_exit(event.clone());
2345    }
2346}
2347
2348fn publish_order_initialized(order: &OrderAny) {
2349    let topic = format!("events.order.{}", order.strategy_id());
2350    let event = OrderEventAny::Initialized(order.init_event().clone());
2351    msgbus::publish_order_event(topic.into(), &event);
2352}
2353
2354fn send_emulator_command(command: TradingCommand) {
2355    log_cmd_send(&command);
2356    let endpoint = MessagingSwitchboard::order_emulator_execute();
2357    msgbus::send_trading_command(endpoint, command);
2358}
2359
2360fn send_algo_command(command: SubmitOrder, exec_algorithm_id: ExecAlgorithmId) {
2361    let id = command.strategy_id;
2362    log::info!("{id} {CMD}{SEND} {command}");
2363
2364    let endpoint = format!("{exec_algorithm_id}.execute");
2365    msgbus::send_any(endpoint.into(), &TradingCommand::SubmitOrder(command));
2366}
2367
2368fn send_risk_command(command: TradingCommand) {
2369    log_cmd_send(&command);
2370    let endpoint = MessagingSwitchboard::risk_engine_queue_execute();
2371    msgbus::send_trading_command(endpoint, command);
2372}
2373
2374fn send_exec_command(command: TradingCommand) {
2375    log_cmd_send(&command);
2376    let endpoint = MessagingSwitchboard::exec_engine_queue_execute();
2377    msgbus::send_trading_command(endpoint, command);
2378}
2379
2380fn log_cmd_send(command: &TradingCommand) {
2381    if let Some(id) = command.strategy_id() {
2382        log::info!("{id} {CMD}{SEND} {command}");
2383    } else {
2384        log::info!("{CMD}{SEND} {command}");
2385    }
2386}
2387
2388fn registered_trader_id(core: &StrategyCore) -> anyhow::Result<TraderId> {
2389    core.trader_id()
2390        .ok_or_else(|| anyhow::anyhow!("Strategy not registered: trader_id is not set"))
2391}
2392
2393fn registered_strategy_id(core: &StrategyCore) -> anyhow::Result<StrategyId> {
2394    core.strategy_id()
2395        .ok_or_else(|| anyhow::anyhow!("Strategy not registered: strategy_id is not set"))
2396}
2397
2398fn required_account_id(order: &OrderAny, operation: &str) -> anyhow::Result<AccountId> {
2399    order.account_id().ok_or_else(|| {
2400        anyhow::anyhow!(
2401            "Cannot generate {operation} event for {}: account_id is not set",
2402            order.client_order_id()
2403        )
2404    })
2405}
2406
2407#[cfg(test)]
2408mod tests {
2409    use std::{cell::RefCell, rc::Rc};
2410
2411    use nautilus_common::{
2412        actor::{
2413            DataActor,
2414            registry::{deregister_actor, try_get_actor_unchecked},
2415        },
2416        cache::{Cache, ORDER_NOT_FOUND},
2417        clock::{Clock, TestClock},
2418        component::{Component, deregister_component, register_component_actor},
2419        enums::ComponentState,
2420        msgbus::{
2421            self, MessagingSwitchboard, TypedHandler, TypedIntoHandler,
2422            stubs::{
2423                TypedIntoMessageSavingHandler, TypedMessageSavingHandler, get_any_saving_handler,
2424                get_typed_into_message_saving_handler, get_typed_message_saving_handler,
2425            },
2426        },
2427        timer::{TimeEvent, TimeEventCallback},
2428    };
2429    use nautilus_core::UnixNanos;
2430    use nautilus_model::{
2431        enums::{
2432            ContingencyType, LiquiditySide, OrderSide, OrderStatus, OrderType,
2433            PositionAdjustmentType, PositionSide, TriggerType,
2434        },
2435        events::{
2436            OrderAccepted, OrderCanceled, OrderFilled, OrderRejected, PositionAdjusted,
2437            order::spec::{
2438                OrderAcceptedSpec, OrderCanceledSpec, OrderEmulatedSpec, OrderExpiredSpec,
2439                OrderFillVoidedSpec, OrderFilledSpec, OrderRejectedSpec,
2440            },
2441        },
2442        identifiers::{
2443            AccountId, ActorId, ClientOrderId, InstrumentId, OrderListId, PositionId, StrategyId,
2444            TradeId, TraderId, VenueOrderId,
2445        },
2446        orderbook::own::OwnOrderBook,
2447        orders::{LimitOrder, MarketOrder, OrderTestBuilder, stubs::TestOrderEventStubs},
2448        stubs::TestDefault,
2449        types::{Currency, Money, Price},
2450    };
2451    use nautilus_portfolio::portfolio::Portfolio;
2452    use rstest::rstest;
2453    use serde_json::Value;
2454
2455    use super::*;
2456    use crate::nautilus_strategy;
2457
2458    #[derive(Debug)]
2459    struct TestStrategy {
2460        core: StrategyCore,
2461        on_order_rejected_called: bool,
2462        on_order_event_called: bool,
2463        on_order_accepted_called: bool,
2464        on_order_canceled_called: bool,
2465        on_order_filled_called: bool,
2466        on_order_fill_voided_called: bool,
2467        on_order_expired_called: bool,
2468        order_event_timeline: Rc<RefCell<Vec<&'static str>>>,
2469        on_position_event_called: bool,
2470        on_position_opened_called: bool,
2471        on_position_changed_called: bool,
2472        on_position_closed_called: bool,
2473    }
2474
2475    #[derive(Debug)]
2476    struct CoreFreeStrategy {
2477        started: bool,
2478    }
2479
2480    #[derive(Debug)]
2481    struct InitializedModifyStrategy {
2482        core: StrategyCore,
2483        modified_quantity: Quantity,
2484    }
2485
2486    #[derive(Debug)]
2487    struct TimerOverrideStrategy {
2488        core: StrategyCore,
2489        gtd_expiries: usize,
2490        market_exit_checks: usize,
2491    }
2492
2493    impl DataActor for CoreFreeStrategy {
2494        fn on_start(&mut self) -> anyhow::Result<()> {
2495            self.started = true;
2496            Ok(())
2497        }
2498    }
2499
2500    impl Strategy for CoreFreeStrategy {}
2501
2502    impl DataActor for InitializedModifyStrategy {}
2503
2504    nautilus_strategy!(InitializedModifyStrategy, {
2505        fn on_order_initialized(&mut self, event: OrderInitialized) {
2506            self.modify_order(
2507                event.client_order_id,
2508                Some(self.modified_quantity),
2509                None,
2510                None,
2511                None,
2512                None,
2513            )
2514            .unwrap();
2515        }
2516    });
2517
2518    impl DataActor for TimerOverrideStrategy {
2519        fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()> {
2520            Strategy::on_time_event(self, event)
2521        }
2522    }
2523
2524    nautilus_strategy!(TimerOverrideStrategy, {
2525        fn check_market_exit(&mut self, _event: TimeEvent) {
2526            self.market_exit_checks += 1;
2527        }
2528
2529        fn expire_gtd_order(&mut self, _event: TimeEvent) {
2530            self.gtd_expiries += 1;
2531        }
2532    });
2533
2534    impl TestStrategy {
2535        fn new(config: StrategyConfig) -> Self {
2536            Self {
2537                core: StrategyCore::new(config),
2538                on_order_rejected_called: false,
2539                on_order_event_called: false,
2540                on_order_accepted_called: false,
2541                on_order_canceled_called: false,
2542                on_order_filled_called: false,
2543                on_order_fill_voided_called: false,
2544                on_order_expired_called: false,
2545                order_event_timeline: Rc::new(RefCell::new(Vec::new())),
2546                on_position_event_called: false,
2547                on_position_opened_called: false,
2548                on_position_changed_called: false,
2549                on_position_closed_called: false,
2550            }
2551        }
2552    }
2553
2554    impl DataActor for TestStrategy {}
2555
2556    nautilus_strategy!(TestStrategy, {
2557        fn on_order_canceled(&mut self, _event: &OrderCanceled) {
2558            self.on_order_canceled_called = true;
2559        }
2560
2561        fn on_order_filled(&mut self, _event: &OrderFilled) {
2562            self.on_order_filled_called = true;
2563            self.order_event_timeline.borrow_mut().push("specific");
2564        }
2565
2566        fn on_order_fill_voided(&mut self, _event: &OrderFillVoided) {
2567            self.on_order_fill_voided_called = true;
2568        }
2569
2570        fn on_order_rejected(&mut self, _event: OrderRejected) {
2571            self.on_order_rejected_called = true;
2572        }
2573
2574        fn on_order_event(&mut self, _event: OrderEventAny) {
2575            self.on_order_event_called = true;
2576            self.order_event_timeline.borrow_mut().push("aggregate");
2577        }
2578
2579        fn on_order_accepted(&mut self, _event: OrderAccepted) {
2580            self.on_order_accepted_called = true;
2581        }
2582
2583        fn on_order_expired(&mut self, _event: OrderExpired) {
2584            self.on_order_expired_called = true;
2585        }
2586
2587        fn on_position_opened(&mut self, _event: PositionOpened) {
2588            self.on_position_opened_called = true;
2589        }
2590
2591        fn on_position_event(&mut self, _event: PositionEvent) {
2592            self.on_position_event_called = true;
2593        }
2594
2595        fn on_position_changed(&mut self, _event: PositionChanged) {
2596            self.on_position_changed_called = true;
2597        }
2598
2599        fn on_position_closed(&mut self, _event: PositionClosed) {
2600            self.on_position_closed_called = true;
2601        }
2602    });
2603
2604    fn create_test_strategy() -> TestStrategy {
2605        let config = StrategyConfig {
2606            strategy_id: Some(StrategyId::from("TEST-001")),
2607            order_id_tag: Some("001".to_string()),
2608            ..Default::default()
2609        };
2610        TestStrategy::new(config)
2611    }
2612
2613    fn register_strategy(strategy: &mut TestStrategy) {
2614        let trader_id = TraderId::from("TRADER-001");
2615        let clock = Rc::new(RefCell::new(TestClock::new()));
2616        let cache = Rc::new(RefCell::new(Cache::default()));
2617        let portfolio = Rc::new(RefCell::new(Portfolio::new(
2618            clock.clone(),
2619            cache.clone(),
2620            None,
2621        )));
2622
2623        strategy
2624            .core
2625            .register(trader_id, clock, cache, portfolio)
2626            .unwrap();
2627        strategy.initialize().unwrap();
2628    }
2629
2630    fn start_strategy(strategy: &mut TestStrategy) {
2631        strategy.start().unwrap();
2632    }
2633
2634    fn stop_strategy(strategy: &mut TestStrategy) {
2635        Component::stop(strategy).unwrap();
2636    }
2637
2638    fn make_filled(client_order_id: ClientOrderId) -> OrderEventAny {
2639        OrderEventAny::Filled(
2640            OrderFilledSpec::builder()
2641                .trader_id(TraderId::from("TRADER-001"))
2642                .strategy_id(StrategyId::from("TEST-001"))
2643                .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2644                .client_order_id(client_order_id)
2645                .venue_order_id(VenueOrderId::test_default())
2646                .account_id(AccountId::from("ACC-001"))
2647                .trade_id(TradeId::test_default())
2648                .last_qty(Quantity::default())
2649                .last_px(Price::default())
2650                .currency(Currency::from("USD"))
2651                .liquidity_side(LiquiditySide::Taker)
2652                .event_id(UUID4::default())
2653                .build(),
2654        )
2655    }
2656
2657    fn make_fill_voided(client_order_id: ClientOrderId, is_reopened: bool) -> OrderEventAny {
2658        OrderEventAny::FillVoided(
2659            OrderFillVoidedSpec::builder()
2660                .trader_id(TraderId::from("TRADER-001"))
2661                .strategy_id(StrategyId::from("TEST-001"))
2662                .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2663                .client_order_id(client_order_id)
2664                .venue_order_id(VenueOrderId::test_default())
2665                .account_id(AccountId::from("ACC-001"))
2666                .is_reopened(is_reopened)
2667                .build(),
2668        )
2669    }
2670
2671    fn make_terminal_fill_voided(client_order_id: ClientOrderId) -> OrderEventAny {
2672        make_fill_voided(client_order_id, false)
2673    }
2674
2675    fn make_canceled(client_order_id: ClientOrderId) -> OrderEventAny {
2676        OrderEventAny::Canceled(
2677            OrderCanceledSpec::builder()
2678                .trader_id(TraderId::from("TRADER-001"))
2679                .strategy_id(StrategyId::from("TEST-001"))
2680                .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2681                .client_order_id(client_order_id)
2682                .account_id(AccountId::from("ACC-001"))
2683                .event_id(UUID4::default())
2684                .build(),
2685        )
2686    }
2687
2688    fn make_rejected(client_order_id: ClientOrderId) -> OrderEventAny {
2689        OrderEventAny::Rejected(
2690            OrderRejectedSpec::builder()
2691                .trader_id(TraderId::from("TRADER-001"))
2692                .strategy_id(StrategyId::from("TEST-001"))
2693                .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2694                .client_order_id(client_order_id)
2695                .account_id(AccountId::from("ACC-001"))
2696                .reason("Test rejection".into())
2697                .event_id(UUID4::default())
2698                .build(),
2699        )
2700    }
2701
2702    fn make_expired(client_order_id: ClientOrderId) -> OrderEventAny {
2703        OrderEventAny::Expired(
2704            OrderExpiredSpec::builder()
2705                .trader_id(TraderId::from("TRADER-001"))
2706                .strategy_id(StrategyId::from("TEST-001"))
2707                .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2708                .client_order_id(client_order_id)
2709                .account_id(AccountId::from("ACC-001"))
2710                .event_id(UUID4::default())
2711                .build(),
2712        )
2713    }
2714
2715    fn make_accepted(client_order_id: ClientOrderId) -> OrderEventAny {
2716        OrderEventAny::Accepted(
2717            OrderAcceptedSpec::builder()
2718                .trader_id(TraderId::from("TRADER-001"))
2719                .strategy_id(StrategyId::from("TEST-001"))
2720                .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2721                .client_order_id(client_order_id)
2722                .venue_order_id(VenueOrderId::test_default())
2723                .account_id(AccountId::from("ACC-001"))
2724                .event_id(UUID4::default())
2725                .build(),
2726        )
2727    }
2728
2729    fn make_accepted_market_order(client_order_id: &str) -> OrderAny {
2730        let mut order = OrderAny::Market(MarketOrder::new(
2731            TraderId::from("TRADER-001"),
2732            StrategyId::from("TEST-001"),
2733            InstrumentId::from("BTCUSDT.BINANCE"),
2734            ClientOrderId::from(client_order_id),
2735            OrderSide::Buy,
2736            Quantity::from(100_000),
2737            TimeInForce::Gtc,
2738            UUID4::new(),
2739            UnixNanos::default(),
2740            false,
2741            false,
2742            None,
2743            None,
2744            None,
2745            None,
2746            None,
2747            None,
2748            None,
2749            None,
2750        ));
2751        let account_id = AccountId::from("ACC-001");
2752        order
2753            .apply(TestOrderEventStubs::submitted(&order, account_id))
2754            .unwrap();
2755        order
2756            .apply(TestOrderEventStubs::accepted(
2757                &order,
2758                account_id,
2759                // Derived per order: a venue order ID has a single owning client order, so
2760                // batch tests holding two accepted orders cannot share the default.
2761                VenueOrderId::from(client_order_id),
2762            ))
2763            .unwrap();
2764        order
2765    }
2766
2767    fn make_accepted_limit_order(client_order_id: &str) -> OrderAny {
2768        let mut order = OrderAny::Limit(LimitOrder::new(
2769            TraderId::from("TRADER-001"),
2770            StrategyId::from("TEST-001"),
2771            InstrumentId::from("BTCUSDT.BINANCE"),
2772            ClientOrderId::from(client_order_id),
2773            OrderSide::Buy,
2774            Quantity::from("1.0"),
2775            Price::from("50000.0"),
2776            TimeInForce::Gtc,
2777            None,
2778            false,
2779            false,
2780            false,
2781            None,
2782            None,
2783            None,
2784            None,
2785            None,
2786            None,
2787            None,
2788            None,
2789            None,
2790            None,
2791            None,
2792            UUID4::new(),
2793            UnixNanos::default(),
2794        ));
2795        let account_id = AccountId::from("ACC-001");
2796        order
2797            .apply(TestOrderEventStubs::submitted(&order, account_id))
2798            .unwrap();
2799        order
2800            .apply(TestOrderEventStubs::accepted(
2801                &order,
2802                account_id,
2803                // Derived per order, as above.
2804                VenueOrderId::from(client_order_id),
2805            ))
2806            .unwrap();
2807        order
2808    }
2809
2810    fn make_initialized_market_order(client_order_id: &str) -> OrderAny {
2811        OrderAny::Market(MarketOrder::new(
2812            TraderId::from("TRADER-001"),
2813            StrategyId::from("TEST-001"),
2814            InstrumentId::from("BTCUSDT.BINANCE"),
2815            ClientOrderId::from(client_order_id),
2816            OrderSide::Buy,
2817            Quantity::from(100_000),
2818            TimeInForce::Gtc,
2819            UUID4::new(),
2820            UnixNanos::default(),
2821            false,
2822            false,
2823            None,
2824            None,
2825            None,
2826            None,
2827            None,
2828            None,
2829            None,
2830            None,
2831        ))
2832    }
2833
2834    fn make_initialized_algorithm_order(client_order_id: &str) -> OrderAny {
2835        OrderAny::Market(MarketOrder::new(
2836            TraderId::from("TRADER-001"),
2837            StrategyId::from("TEST-001"),
2838            InstrumentId::from("BTCUSDT.BINANCE"),
2839            ClientOrderId::from(client_order_id),
2840            OrderSide::Buy,
2841            Quantity::from(100_000),
2842            TimeInForce::Gtc,
2843            UUID4::new(),
2844            UnixNanos::default(),
2845            false,
2846            false,
2847            None,
2848            None,
2849            None,
2850            None,
2851            Some(ExecAlgorithmId::from("TWAP")),
2852            None,
2853            Some(ClientOrderId::from(client_order_id)),
2854            None,
2855        ))
2856    }
2857
2858    fn add_order_to_cache(strategy: &TestStrategy, order: &OrderAny) {
2859        let cache_rc = strategy.core.cache_rc();
2860        let mut cache = cache_rc.borrow_mut();
2861        cache.add_order(order.clone(), None, None, true).unwrap();
2862    }
2863
2864    fn make_submit_command(order: &OrderAny) -> SubmitOrder {
2865        SubmitOrder::new(
2866            order.trader_id(),
2867            None,
2868            order.strategy_id(),
2869            order.instrument_id(),
2870            order.client_order_id(),
2871            order.init_event().clone(),
2872            order.exec_algorithm_id(),
2873            None,
2874            None,
2875            UUID4::new(),
2876            UnixNanos::default(),
2877            None, // correlation_id
2878        )
2879    }
2880
2881    fn add_order_to_cache_and_own_book(strategy: &TestStrategy, order: &OrderAny) {
2882        let cache_rc = strategy.core.cache_rc();
2883        let mut cache = cache_rc.borrow_mut();
2884        cache.add_order(order.clone(), None, None, true).unwrap();
2885        cache
2886            .add_own_order_book(OwnOrderBook::new(order.instrument_id()))
2887            .unwrap();
2888        cache.update_own_order_book(order);
2889    }
2890
2891    fn make_position_opened() -> PositionEvent {
2892        PositionEvent::PositionOpened(PositionOpened {
2893            trader_id: TraderId::from("TRADER-001"),
2894            strategy_id: StrategyId::from("TEST-001"),
2895            instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
2896            position_id: PositionId::test_default(),
2897            account_id: AccountId::from("ACC-001"),
2898            opening_order_id: ClientOrderId::from("O-001"),
2899            entry: OrderSide::Buy,
2900            side: PositionSide::Long,
2901            signed_qty: 1.0,
2902            quantity: Quantity::default(),
2903            last_qty: Quantity::default(),
2904            last_px: Price::default(),
2905            currency: Currency::from("USD"),
2906            avg_px_open: 0.0,
2907            realized_pnl: None,
2908            event_id: UUID4::default(),
2909            ts_event: UnixNanos::default(),
2910            ts_init: UnixNanos::default(),
2911        })
2912    }
2913
2914    fn make_position_changed() -> PositionEvent {
2915        let currency = Currency::from("USD");
2916        PositionEvent::PositionChanged(PositionChanged {
2917            trader_id: TraderId::from("TRADER-001"),
2918            strategy_id: StrategyId::from("TEST-001"),
2919            instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
2920            position_id: PositionId::test_default(),
2921            account_id: AccountId::from("ACC-001"),
2922            opening_order_id: ClientOrderId::from("O-001"),
2923            entry: OrderSide::Buy,
2924            side: PositionSide::Long,
2925            signed_qty: 2.0,
2926            quantity: Quantity::default(),
2927            peak_quantity: Quantity::default(),
2928            last_qty: Quantity::default(),
2929            last_px: Price::default(),
2930            currency,
2931            avg_px_open: 0.0,
2932            avg_px_close: None,
2933            realized_return: 0.0,
2934            realized_pnl: None,
2935            unrealized_pnl: Money::zero(currency),
2936            event_id: UUID4::default(),
2937            ts_opened: UnixNanos::default(),
2938            ts_event: UnixNanos::default(),
2939            ts_init: UnixNanos::default(),
2940        })
2941    }
2942
2943    fn make_position_closed() -> PositionEvent {
2944        let currency = Currency::from("USD");
2945        PositionEvent::PositionClosed(PositionClosed {
2946            trader_id: TraderId::from("TRADER-001"),
2947            strategy_id: StrategyId::from("TEST-001"),
2948            instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
2949            position_id: PositionId::test_default(),
2950            account_id: AccountId::from("ACC-001"),
2951            opening_order_id: ClientOrderId::from("O-001"),
2952            closing_order_id: Some(ClientOrderId::from("O-002")),
2953            entry: OrderSide::Buy,
2954            side: PositionSide::Flat,
2955            signed_qty: 0.0,
2956            quantity: Quantity::default(),
2957            peak_quantity: Quantity::default(),
2958            last_qty: Quantity::default(),
2959            last_px: Price::default(),
2960            currency,
2961            avg_px_open: 0.0,
2962            avg_px_close: None,
2963            realized_return: 0.0,
2964            realized_pnl: None,
2965            unrealized_pnl: Money::zero(currency),
2966            duration: 0,
2967            event_id: UUID4::default(),
2968            ts_opened: UnixNanos::default(),
2969            ts_closed: None,
2970            ts_event: UnixNanos::default(),
2971            ts_init: UnixNanos::default(),
2972        })
2973    }
2974
2975    fn make_position_adjusted() -> PositionEvent {
2976        PositionEvent::PositionAdjusted(PositionAdjusted {
2977            trader_id: TraderId::from("TRADER-001"),
2978            strategy_id: StrategyId::from("TEST-001"),
2979            instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
2980            position_id: PositionId::test_default(),
2981            account_id: AccountId::from("ACC-001"),
2982            adjustment_type: PositionAdjustmentType::Funding,
2983            quantity_change: None,
2984            pnl_change: None,
2985            reason: None,
2986            event_id: UUID4::default(),
2987            ts_event: UnixNanos::default(),
2988            ts_init: UnixNanos::default(),
2989        })
2990    }
2991
2992    #[rstest]
2993    fn test_strategy_creation() {
2994        let strategy = create_test_strategy();
2995        assert_eq!(strategy.strategy_id(), Some(StrategyId::from("TEST-001")));
2996        assert!(!strategy.on_order_rejected_called);
2997        assert!(!strategy.on_position_opened_called);
2998    }
2999
3000    #[rstest]
3001    fn test_strategy_registration() {
3002        let mut strategy = create_test_strategy();
3003        register_strategy(&mut strategy);
3004
3005        assert!(strategy.is_registered());
3006        let _ = strategy.order().generate_client_order_id();
3007        let _ = strategy.portfolio().is_initialized();
3008    }
3009
3010    #[rstest]
3011    fn test_strategy_native_methods_are_available_on_strategy_type() {
3012        let mut strategy = create_test_strategy();
3013        register_strategy(&mut strategy);
3014
3015        drop(strategy.order_factory());
3016
3017        assert!(Rc::ptr_eq(
3018            &strategy.order_factory_rc(),
3019            strategy.core.order_factory.as_ref().unwrap()
3020        ));
3021        assert!(Rc::ptr_eq(
3022            &strategy.portfolio_rc(),
3023            strategy.core.portfolio.as_ref().unwrap()
3024        ));
3025    }
3026
3027    #[rstest]
3028    fn test_handle_order_event_dispatches_to_handler() {
3029        let mut strategy = create_test_strategy();
3030        register_strategy(&mut strategy);
3031        start_strategy(&mut strategy);
3032
3033        let event = make_rejected(ClientOrderId::from("O-001"));
3034
3035        strategy.handle_order_event(event);
3036
3037        assert!(strategy.on_order_rejected_called);
3038        assert!(strategy.on_order_event_called);
3039    }
3040
3041    #[rstest]
3042    fn test_dispatch_manager_actions_routes_every_action() {
3043        let mut strategy = create_test_strategy();
3044        register_strategy(&mut strategy);
3045        let strategy_id = StrategyId::from("TEST-001");
3046        let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
3047        let initialized_order = make_initialized_market_order("O-MANAGER-INIT");
3048        let emulator_order = OrderTestBuilder::new(OrderType::StopMarket)
3049            .trader_id(TraderId::from("TRADER-001"))
3050            .strategy_id(strategy_id)
3051            .instrument_id(instrument_id)
3052            .client_order_id(ClientOrderId::from("O-MANAGER-EMULATOR"))
3053            .side(OrderSide::Buy)
3054            .trigger_price(Price::from("51000.0"))
3055            .quantity(Quantity::from(100_000))
3056            .emulation_trigger(TriggerType::BidAsk)
3057            .build();
3058        let risk_order = make_initialized_market_order("O-MANAGER-RISK");
3059        let algorithm_order = make_initialized_algorithm_order("O-MANAGER-ALGORITHM");
3060        let cancel_order = OrderTestBuilder::new(OrderType::Limit)
3061            .trader_id(TraderId::from("TRADER-001"))
3062            .strategy_id(strategy_id)
3063            .instrument_id(instrument_id)
3064            .client_order_id(ClientOrderId::from("O-MANAGER-CANCEL"))
3065            .side(OrderSide::Buy)
3066            .price(Price::from("50000.0"))
3067            .quantity(Quantity::from(100_000))
3068            .submit(true)
3069            .build();
3070        let missing_cancel_order = OrderTestBuilder::new(OrderType::Limit)
3071            .trader_id(TraderId::from("TRADER-001"))
3072            .strategy_id(strategy_id)
3073            .instrument_id(instrument_id)
3074            .client_order_id(ClientOrderId::from("O-MANAGER-MISSING-CANCEL"))
3075            .side(OrderSide::Buy)
3076            .price(Price::from("50000.0"))
3077            .quantity(Quantity::from(100_000))
3078            .submit(true)
3079            .build();
3080        let modify_order = OrderTestBuilder::new(OrderType::Limit)
3081            .trader_id(TraderId::from("TRADER-001"))
3082            .strategy_id(strategy_id)
3083            .instrument_id(instrument_id)
3084            .client_order_id(ClientOrderId::from("O-MANAGER-MODIFY"))
3085            .side(OrderSide::Buy)
3086            .price(Price::from("50000.0"))
3087            .quantity(Quantity::from(100_000))
3088            .submit(true)
3089            .build();
3090        add_order_to_cache(&strategy, &cancel_order);
3091        add_order_to_cache(&strategy, &modify_order);
3092
3093        let (order_handler, order_events) =
3094            get_typed_message_saving_handler(Some(Ustr::from("manager-action-order-events")));
3095        let order_topic = format!("events.order.{strategy_id}");
3096        msgbus::subscribe_order_events(order_topic.clone().into(), order_handler.clone(), None);
3097        let (emulator_handler, emulator_messages): (
3098            _,
3099            TypedIntoMessageSavingHandler<TradingCommand>,
3100        ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
3101        msgbus::register_trading_command_endpoint(
3102            MessagingSwitchboard::order_emulator_execute(),
3103            emulator_handler,
3104        );
3105        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3106            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3107        msgbus::register_trading_command_endpoint(
3108            MessagingSwitchboard::risk_engine_queue_execute(),
3109            risk_handler,
3110        );
3111        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3112            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3113        msgbus::register_trading_command_endpoint(
3114            MessagingSwitchboard::exec_engine_queue_execute(),
3115            exec_handler,
3116        );
3117        let (algorithm_handler, algorithm_messages) =
3118            get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
3119        msgbus::register_any("TWAP.execute".into(), algorithm_handler);
3120
3121        strategy.dispatch_manager_actions(vec![
3122            OrderManagerAction::CancelLocal(missing_cancel_order),
3123            OrderManagerAction::PublishInitialized(OrderEventAny::Initialized(
3124                initialized_order.init_event().clone(),
3125            )),
3126            OrderManagerAction::SubmitToEmulator(make_submit_command(&emulator_order)),
3127            OrderManagerAction::SubmitToRisk(make_submit_command(&risk_order)),
3128            OrderManagerAction::SubmitToAlgorithm {
3129                command: make_submit_command(&algorithm_order),
3130                exec_algorithm_id: ExecAlgorithmId::from("TWAP"),
3131            },
3132            OrderManagerAction::CancelLocal(cancel_order.clone()),
3133            OrderManagerAction::ModifyLocalQuantity {
3134                order: modify_order.clone(),
3135                quantity: Quantity::from(50_000),
3136            },
3137            OrderManagerAction::ModifyLocalQuantity {
3138                order: modify_order.clone(),
3139                quantity: Quantity::from(100_000),
3140            },
3141        ]);
3142        msgbus::unsubscribe_order_events(order_topic.into(), &order_handler);
3143
3144        let order_events = order_events.get_messages();
3145        assert_eq!(order_events.len(), 3);
3146        assert!(matches!(
3147            &order_events[0],
3148            OrderEventAny::Initialized(event)
3149                if event.client_order_id == initialized_order.client_order_id()
3150        ));
3151        assert!(matches!(
3152            &order_events[1],
3153            OrderEventAny::PendingCancel(event)
3154                if event.client_order_id == cancel_order.client_order_id()
3155        ));
3156        assert!(matches!(
3157            &order_events[2],
3158            OrderEventAny::PendingUpdate(event)
3159                if event.client_order_id == modify_order.client_order_id()
3160        ));
3161        assert!(matches!(
3162            emulator_messages.get_messages().as_slice(),
3163            [TradingCommand::SubmitOrder(command)]
3164                if command.client_order_id == emulator_order.client_order_id()
3165        ));
3166        assert!(matches!(
3167            risk_messages.get_messages().as_slice(),
3168            [
3169                TradingCommand::SubmitOrder(submit),
3170                TradingCommand::ModifyOrder(first_modify),
3171                TradingCommand::ModifyOrder(second_modify),
3172            ]
3173
3174                if submit.client_order_id == risk_order.client_order_id()
3175                    && first_modify.client_order_id == modify_order.client_order_id()
3176                    && first_modify.quantity == Some(Quantity::from(50_000))
3177                    && second_modify.client_order_id == modify_order.client_order_id()
3178                    && second_modify.quantity == Some(Quantity::from(100_000))
3179        ));
3180        assert!(matches!(
3181            algorithm_messages.get_messages().as_slice(),
3182            [TradingCommand::SubmitOrder(command)]
3183                if command.client_order_id == algorithm_order.client_order_id()
3184        ));
3185        assert!(matches!(
3186            exec_messages.get_messages().as_slice(),
3187            [TradingCommand::CancelOrder(command)]
3188                if command.client_order_id == cancel_order.client_order_id()
3189        ));
3190    }
3191
3192    #[rstest]
3193    #[case::disabled(false)]
3194    #[case::enabled(true)]
3195    fn test_contingent_manager_dispatch_precedes_user_handlers_and_is_idempotent(
3196        #[case] manage_contingent_orders: bool,
3197    ) {
3198        let mut strategy = TestStrategy::new(StrategyConfig {
3199            strategy_id: Some(StrategyId::from("TEST-001")),
3200            order_id_tag: Some("001".to_string()),
3201            manage_contingent_orders,
3202            ..Default::default()
3203        });
3204        register_strategy(&mut strategy);
3205        start_strategy(&mut strategy);
3206        let timeline = strategy.order_event_timeline.clone();
3207        let endpoint_timeline = timeline.clone();
3208        let exec_handler = TypedIntoHandler::from(move |command: TradingCommand| {
3209            assert!(matches!(command, TradingCommand::CancelOrder(_)));
3210            endpoint_timeline.borrow_mut().push("endpoint");
3211        });
3212        msgbus::register_trading_command_endpoint(
3213            MessagingSwitchboard::exec_engine_queue_execute(),
3214            exec_handler,
3215        );
3216        let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
3217        let parent_id = ClientOrderId::from("O-CONTINGENT-PARENT");
3218        let sibling_id = ClientOrderId::from("O-CONTINGENT-SIBLING");
3219        let parent = OrderTestBuilder::new(OrderType::Limit)
3220            .trader_id(TraderId::from("TRADER-001"))
3221            .strategy_id(StrategyId::from("TEST-001"))
3222            .instrument_id(instrument_id)
3223            .client_order_id(parent_id)
3224            .side(OrderSide::Buy)
3225            .price(Price::from("50000.0"))
3226            .quantity(Quantity::from(100_000))
3227            .contingency_type(ContingencyType::Oco)
3228            .linked_order_ids(vec![parent_id, sibling_id])
3229            .submit(true)
3230            .build();
3231        let sibling = OrderTestBuilder::new(OrderType::Limit)
3232            .trader_id(TraderId::from("TRADER-001"))
3233            .strategy_id(StrategyId::from("TEST-001"))
3234            .instrument_id(instrument_id)
3235            .client_order_id(sibling_id)
3236            .side(OrderSide::Sell)
3237            .price(Price::from("51000.0"))
3238            .quantity(Quantity::from(100_000))
3239            .submit(true)
3240            .build();
3241        add_order_to_cache(&strategy, &parent);
3242        add_order_to_cache(&strategy, &sibling);
3243        let event = OrderEventAny::Filled(
3244            OrderFilledSpec::builder()
3245                .trader_id(parent.trader_id())
3246                .strategy_id(parent.strategy_id())
3247                .instrument_id(parent.instrument_id())
3248                .client_order_id(parent.client_order_id())
3249                .venue_order_id(VenueOrderId::from("V-CONTINGENT-PARENT"))
3250                .account_id(AccountId::from("ACCOUNT-001"))
3251                .trade_id(TradeId::from("T-CONTINGENT-PARENT"))
3252                .order_side(parent.order_side())
3253                .order_type(parent.order_type())
3254                .last_qty(parent.quantity())
3255                .last_px(Price::from("50000.0"))
3256                .liquidity_side(LiquiditySide::Taker)
3257                .build(),
3258        );
3259        strategy
3260            .core
3261            .cache_rc()
3262            .borrow_mut()
3263            .update_order(&event)
3264            .unwrap();
3265
3266        strategy.handle_order_event(event.clone());
3267        strategy.handle_order_event(event);
3268
3269        let timeline = timeline.borrow();
3270        let sibling_status = strategy
3271            .core
3272            .cache_ref()
3273            .order(&sibling_id)
3274            .unwrap()
3275            .status();
3276
3277        if manage_contingent_orders {
3278            assert_eq!(
3279                timeline.as_slice(),
3280                ["endpoint", "specific", "aggregate", "specific", "aggregate"]
3281            );
3282            assert_eq!(sibling_status, OrderStatus::PendingCancel);
3283        } else {
3284            assert_eq!(
3285                timeline.as_slice(),
3286                ["specific", "aggregate", "specific", "aggregate"]
3287            );
3288            assert_eq!(sibling_status, OrderStatus::Submitted);
3289        }
3290    }
3291
3292    #[rstest]
3293    fn test_handle_order_fill_voided_dispatches_to_specific_handler() {
3294        let mut strategy = create_test_strategy();
3295        register_strategy(&mut strategy);
3296        start_strategy(&mut strategy);
3297
3298        strategy.handle_order_event(make_fill_voided(ClientOrderId::from("O-001"), false));
3299
3300        assert!(strategy.on_order_fill_voided_called);
3301        assert!(strategy.on_order_event_called);
3302    }
3303
3304    #[rstest]
3305    #[case::opened(make_position_opened())]
3306    #[case::changed(make_position_changed())]
3307    #[case::closed(make_position_closed())]
3308    fn test_handle_position_event_dispatches_to_handler(#[case] event: PositionEvent) {
3309        let mut strategy = create_test_strategy();
3310        register_strategy(&mut strategy);
3311        start_strategy(&mut strategy);
3312
3313        let expected_opened = matches!(event, PositionEvent::PositionOpened(_));
3314        let expected_changed = matches!(event, PositionEvent::PositionChanged(_));
3315        let expected_closed = matches!(event, PositionEvent::PositionClosed(_));
3316
3317        strategy.handle_position_event(event);
3318
3319        assert_eq!(strategy.on_position_opened_called, expected_opened);
3320        assert_eq!(strategy.on_position_changed_called, expected_changed);
3321        assert_eq!(strategy.on_position_closed_called, expected_closed);
3322        assert!(strategy.on_position_event_called);
3323    }
3324
3325    #[rstest]
3326    fn test_handle_position_event_skips_dispatch_when_stopped() {
3327        let mut strategy = create_test_strategy();
3328        register_strategy(&mut strategy);
3329        start_strategy(&mut strategy);
3330        stop_strategy(&mut strategy);
3331        assert_eq!(strategy.state(), ComponentState::Stopped);
3332
3333        strategy.handle_position_event(make_position_opened());
3334
3335        assert!(!strategy.on_position_event_called);
3336        assert!(!strategy.on_position_opened_called);
3337    }
3338
3339    #[rstest]
3340    fn test_handle_position_event_skips_dispatch_for_adjusted() {
3341        let mut strategy = create_test_strategy();
3342        register_strategy(&mut strategy);
3343        start_strategy(&mut strategy);
3344
3345        strategy.handle_position_event(make_position_adjusted());
3346
3347        assert!(!strategy.on_position_event_called);
3348        assert!(!strategy.on_position_opened_called);
3349        assert!(!strategy.on_position_changed_called);
3350        assert!(!strategy.on_position_closed_called);
3351    }
3352
3353    #[rstest]
3354    fn test_strategy_default_handlers_do_not_panic() {
3355        let mut strategy = create_test_strategy();
3356
3357        strategy.on_order_initialized(OrderInitialized::default());
3358        strategy.on_order_event(OrderEventAny::Accepted(OrderAccepted::default()));
3359        strategy.on_order_denied(OrderDenied::default());
3360        strategy.on_order_emulated(OrderEmulated::default());
3361        strategy.on_order_released(OrderReleased::default());
3362        strategy.on_order_submitted(OrderSubmitted::default());
3363        strategy.on_order_rejected(OrderRejected::default());
3364        strategy.on_order_canceled(&OrderCanceled::default());
3365        strategy.on_order_expired(OrderExpired::default());
3366        strategy.on_order_triggered(OrderTriggered::default());
3367        strategy.on_order_pending_update(OrderPendingUpdate::default());
3368        strategy.on_order_pending_cancel(OrderPendingCancel::default());
3369        strategy.on_order_modify_rejected(OrderModifyRejected::default());
3370        strategy.on_order_cancel_rejected(OrderCancelRejected::default());
3371        strategy.on_order_updated(OrderUpdated::default());
3372        strategy.on_order_filled(&OrderFilledSpec::builder().build());
3373        strategy.on_order_fill_voided(&OrderFillVoidedSpec::builder().build());
3374        strategy.on_position_event(make_position_opened());
3375    }
3376
3377    #[rstest]
3378    fn test_submit_order_publishes_order_initialized_after_cache_insert_before_send() {
3379        let mut strategy = create_test_strategy();
3380        register_strategy(&mut strategy);
3381
3382        let order = make_initialized_market_order("O-20250208-INIT-001");
3383        let client_order_id = order.client_order_id();
3384        let cache_rc = strategy.core.cache_rc();
3385        let timeline = Rc::new(RefCell::new(Vec::new()));
3386        let event_messages = Rc::new(RefCell::new(Vec::new()));
3387
3388        let event_handler = {
3389            let event_messages = event_messages.clone();
3390            let timeline = timeline.clone();
3391            TypedHandler::from_with_id("events.order.initialized", move |event: &OrderEventAny| {
3392                assert!(cache_rc.borrow().order_exists(&client_order_id));
3393                assert!(matches!(event, OrderEventAny::Initialized(_)));
3394                event_messages.borrow_mut().push(event.clone());
3395                timeline.borrow_mut().push("init");
3396            })
3397        };
3398        let risk_handler = {
3399            let timeline = timeline.clone();
3400            TypedIntoHandler::from_with_id(
3401                "RiskEngine.queue_execute",
3402                move |command: TradingCommand| {
3403                    assert!(matches!(command, TradingCommand::SubmitOrder(_)));
3404                    timeline.borrow_mut().push("command");
3405                },
3406            )
3407        };
3408        msgbus::register_trading_command_endpoint(
3409            MessagingSwitchboard::risk_engine_queue_execute(),
3410            risk_handler,
3411        );
3412
3413        let topic = format!("events.order.{}", order.strategy_id());
3414        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3415
3416        strategy
3417            .submit_order(order.clone(), None, None, None)
3418            .unwrap();
3419
3420        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3421
3422        let event_messages = event_messages.borrow();
3423        assert_eq!(event_messages.len(), 1);
3424        assert_eq!(
3425            event_messages[0],
3426            OrderEventAny::Initialized(order.init_event().clone())
3427        );
3428        assert_eq!(timeline.borrow().as_slice(), &["init", "command"]);
3429    }
3430
3431    #[rstest]
3432    fn test_submit_order_routes_emulated_order_to_order_emulator() {
3433        let mut strategy = create_test_strategy();
3434        register_strategy(&mut strategy);
3435        let (emulator_handler, emulator_messages): (
3436            _,
3437            TypedIntoMessageSavingHandler<TradingCommand>,
3438        ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
3439        msgbus::register_trading_command_endpoint(
3440            MessagingSwitchboard::order_emulator_execute(),
3441            emulator_handler,
3442        );
3443        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3444            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3445        msgbus::register_trading_command_endpoint(
3446            MessagingSwitchboard::risk_engine_queue_execute(),
3447            risk_handler,
3448        );
3449        let order = OrderTestBuilder::new(OrderType::StopMarket)
3450            .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
3451            .client_order_id(ClientOrderId::from("O-20250208-EMULATED-001"))
3452            .side(OrderSide::Buy)
3453            .trigger_price(Price::from("51000.0"))
3454            .quantity(Quantity::from(100_000))
3455            .emulation_trigger(TriggerType::BidAsk)
3456            .build();
3457        let client_order_id = order.client_order_id();
3458
3459        strategy.submit_order(order, None, None, None).unwrap();
3460
3461        let emulator_messages = emulator_messages.get_messages();
3462        assert_eq!(emulator_messages.len(), 1);
3463        assert!(matches!(
3464            emulator_messages.first(),
3465            Some(TradingCommand::SubmitOrder(command))
3466                if command.client_order_id == client_order_id
3467        ));
3468        assert!(risk_messages.get_messages().is_empty());
3469    }
3470
3471    #[rstest]
3472    fn test_submit_order_errors_when_strategy_not_registered() {
3473        let mut strategy = create_test_strategy();
3474        let order = make_initialized_market_order("O-20250208-UNREGISTERED-001");
3475
3476        let err = strategy
3477            .submit_order(order, None, None, None)
3478            .unwrap_err()
3479            .to_string();
3480
3481        assert_eq!(err, "Strategy not registered: trader_id is not set");
3482    }
3483
3484    #[rstest]
3485    fn test_submit_order_uses_stored_strategy_id_when_actor_id_diverges() {
3486        let mut strategy = create_test_strategy();
3487        register_strategy(&mut strategy);
3488        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3489            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3490        msgbus::register_trading_command_endpoint(
3491            MessagingSwitchboard::risk_engine_queue_execute(),
3492            risk_handler,
3493        );
3494
3495        // A hyphenless actor ID has no strategy ID form, so any re-parse of it panics
3496        strategy.core.actor.actor_id = ActorId::from("Strategy");
3497        let order = make_initialized_market_order("O-20250208-DIVERGED-001");
3498
3499        strategy.submit_order(order, None, None, None).unwrap();
3500
3501        let risk_messages = risk_messages.get_messages();
3502        assert_eq!(risk_messages.len(), 1);
3503        let Some(TradingCommand::SubmitOrder(command)) = risk_messages.first() else {
3504            panic!("Expected a SubmitOrder command, was {risk_messages:?}");
3505        };
3506        assert_eq!(command.strategy_id, StrategyId::from("TEST-001"));
3507        assert_eq!(
3508            command.client_order_id,
3509            ClientOrderId::from("O-20250208-DIVERGED-001")
3510        );
3511    }
3512
3513    #[rstest]
3514    fn test_required_account_id_errors_when_missing_for_strategy_event() {
3515        let order = make_initialized_market_order("O-20250208-NO-ACCOUNT-001");
3516
3517        let err = required_account_id(&order, "pending cancel")
3518            .unwrap_err()
3519            .to_string();
3520
3521        assert_eq!(
3522            err,
3523            "Cannot generate pending cancel event for O-20250208-NO-ACCOUNT-001: \
3524             account_id is not set"
3525        );
3526    }
3527
3528    #[rstest]
3529    fn test_submit_order_rejects_non_initialized_without_events() {
3530        let mut strategy = create_test_strategy();
3531        register_strategy(&mut strategy);
3532
3533        let order = make_accepted_market_order("O-20250208-ACCEPTED-001");
3534        let topic = format!("events.order.{}", order.strategy_id());
3535        let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
3536            get_typed_message_saving_handler(Some(Ustr::from("events.order.invalid")));
3537
3538        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3539        let result = strategy.submit_order(order, None, None, None);
3540
3541        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3542
3543        assert!(result.is_err());
3544        assert!(
3545            result
3546                .unwrap_err()
3547                .to_string()
3548                .contains("expected INITIALIZED")
3549        );
3550        assert!(event_messages.get_messages().is_empty());
3551    }
3552
3553    #[rstest]
3554    fn test_submit_order_returns_error_when_cache_already_borrowed() {
3555        let mut strategy = create_test_strategy();
3556        register_strategy(&mut strategy);
3557
3558        let order = make_initialized_market_order("O-20250208-BORROWED-001");
3559        let cache_rc = strategy.core.cache_rc();
3560        let _cache = cache_rc.borrow();
3561
3562        let result = catch_unwind(AssertUnwindSafe(|| {
3563            strategy.submit_order(order, None, None, None)
3564        }));
3565
3566        let err = result
3567            .expect("submit_order should not panic")
3568            .unwrap_err()
3569            .to_string();
3570
3571        assert_eq!(
3572            err,
3573            "Cannot submit order O-20250208-BORROWED-001: cache is currently borrowed"
3574        );
3575    }
3576
3577    #[rstest]
3578    fn test_submit_order_list_publishes_order_initialized_after_cache_insert_before_send() {
3579        let mut strategy = create_test_strategy();
3580        register_strategy(&mut strategy);
3581
3582        let order_list_id = OrderListId::from("OL-20250208-LIST-INIT");
3583        let mut orders = vec![
3584            make_initialized_market_order("O-20250208-LIST-INIT-001"),
3585            make_initialized_market_order("O-20250208-LIST-INIT-002"),
3586        ];
3587
3588        for order in &mut orders {
3589            order.set_order_list_id(order_list_id);
3590        }
3591
3592        let client_order_id1 = orders[0].client_order_id();
3593        let client_order_id2 = orders[1].client_order_id();
3594        let cache_rc = strategy.core.cache_rc();
3595        let timeline = Rc::new(RefCell::new(Vec::new()));
3596        let event_messages = Rc::new(RefCell::new(Vec::new()));
3597
3598        let event_handler = {
3599            let event_messages = event_messages.clone();
3600            let timeline = timeline.clone();
3601            TypedHandler::from_with_id(
3602                "events.order.list_initialized",
3603                move |event: &OrderEventAny| {
3604                    match event {
3605                        OrderEventAny::Initialized(e) if e.client_order_id == client_order_id1 => {
3606                            let cache = cache_rc.borrow();
3607                            assert!(cache.order_exists(&client_order_id1));
3608                            assert!(cache.order_exists(&client_order_id2));
3609                            assert!(cache.order_list_exists(&order_list_id));
3610                            let order_list = cache.order_list(&order_list_id).unwrap();
3611                            assert_eq!(
3612                                order_list.client_order_ids.as_slice(),
3613                                &[client_order_id1, client_order_id2]
3614                            );
3615                            timeline.borrow_mut().push("init1");
3616                        }
3617                        OrderEventAny::Initialized(e) if e.client_order_id == client_order_id2 => {
3618                            assert!(cache_rc.borrow().order_exists(&client_order_id2));
3619                            timeline.borrow_mut().push("init2");
3620                        }
3621                        _ => panic!("unexpected order event {event:?}"),
3622                    }
3623                    event_messages.borrow_mut().push(event.clone());
3624                },
3625            )
3626        };
3627        let risk_handler = {
3628            let timeline = timeline.clone();
3629            TypedIntoHandler::from_with_id(
3630                "RiskEngine.queue_execute",
3631                move |command: TradingCommand| {
3632                    assert!(matches!(command, TradingCommand::SubmitOrderList(_)));
3633                    timeline.borrow_mut().push("command");
3634                },
3635            )
3636        };
3637        msgbus::register_trading_command_endpoint(
3638            MessagingSwitchboard::risk_engine_queue_execute(),
3639            risk_handler,
3640        );
3641
3642        let topic = format!("events.order.{}", orders[0].strategy_id());
3643        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3644
3645        strategy
3646            .submit_order_list(orders.clone(), None, None, None)
3647            .unwrap();
3648
3649        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3650
3651        let event_messages = event_messages.borrow();
3652        assert_eq!(event_messages.len(), 2);
3653        assert_eq!(
3654            event_messages[0],
3655            OrderEventAny::Initialized(orders[0].init_event().clone())
3656        );
3657        assert_eq!(
3658            event_messages[1],
3659            OrderEventAny::Initialized(orders[1].init_event().clone())
3660        );
3661        assert_eq!(timeline.borrow().as_slice(), &["init1", "init2", "command"]);
3662    }
3663
3664    #[rstest]
3665    fn test_submit_order_list_returns_error_when_cache_already_borrowed() {
3666        let mut strategy = create_test_strategy();
3667        register_strategy(&mut strategy);
3668
3669        let order_list_id = OrderListId::from("OL-20250208-BORROWED");
3670        let mut orders = vec![
3671            make_initialized_market_order("O-20250208-LIST-BORROWED-001"),
3672            make_initialized_market_order("O-20250208-LIST-BORROWED-002"),
3673        ];
3674
3675        for order in &mut orders {
3676            order.set_order_list_id(order_list_id);
3677        }
3678
3679        let cache_rc = strategy.core.cache_rc();
3680        let _cache = cache_rc.borrow();
3681
3682        let result = catch_unwind(AssertUnwindSafe(|| {
3683            strategy.submit_order_list(orders, None, None, None)
3684        }));
3685
3686        let err = result
3687            .expect("submit_order_list should not panic")
3688            .unwrap_err()
3689            .to_string();
3690
3691        assert_eq!(
3692            err,
3693            "Cannot submit order list OL-20250208-BORROWED: cache is currently borrowed"
3694        );
3695    }
3696
3697    #[rstest]
3698    fn test_submit_order_list_create_list_branch_publishes_init_after_cache_insert() {
3699        let mut strategy = create_test_strategy();
3700        register_strategy(&mut strategy);
3701
3702        let orders = vec![
3703            make_initialized_market_order("O-20250208-LIST-CREATE-001"),
3704            make_initialized_market_order("O-20250208-LIST-CREATE-002"),
3705        ];
3706
3707        let client_order_id1 = orders[0].client_order_id();
3708        let client_order_id2 = orders[1].client_order_id();
3709        let cache_rc = strategy.core.cache_rc();
3710        let timeline = Rc::new(RefCell::new(Vec::new()));
3711        let event_messages = Rc::new(RefCell::new(Vec::new()));
3712
3713        let event_handler = {
3714            let event_messages = event_messages.clone();
3715            let timeline = timeline.clone();
3716            TypedHandler::from_with_id(
3717                "events.order.list_create_initialized",
3718                move |event: &OrderEventAny| {
3719                    match event {
3720                        OrderEventAny::Initialized(e) if e.client_order_id == client_order_id1 => {
3721                            let cache = cache_rc.borrow();
3722                            let cached_order1 = cache.order(&client_order_id1).unwrap();
3723                            let cached_order2 = cache.order(&client_order_id2).unwrap();
3724                            let order_list_id = cached_order1.order_list_id().unwrap();
3725                            assert_eq!(cached_order2.order_list_id(), Some(order_list_id));
3726                            assert_eq!(e.order_list_id, Some(order_list_id));
3727                            assert!(cache.order_list_exists(&order_list_id));
3728                            let order_list = cache.order_list(&order_list_id).unwrap();
3729                            assert_eq!(
3730                                order_list.client_order_ids.as_slice(),
3731                                &[client_order_id1, client_order_id2]
3732                            );
3733                            timeline.borrow_mut().push("init1");
3734                        }
3735                        OrderEventAny::Initialized(e) if e.client_order_id == client_order_id2 => {
3736                            let cache = cache_rc.borrow();
3737                            let cached_order = cache.order(&client_order_id2).unwrap();
3738                            assert_eq!(e.order_list_id, cached_order.order_list_id());
3739                            timeline.borrow_mut().push("init2");
3740                        }
3741                        _ => panic!("unexpected order event {event:?}"),
3742                    }
3743                    event_messages.borrow_mut().push(event.clone());
3744                },
3745            )
3746        };
3747        let risk_handler = {
3748            let timeline = timeline.clone();
3749            TypedIntoHandler::from_with_id(
3750                "RiskEngine.queue_execute",
3751                move |command: TradingCommand| {
3752                    let TradingCommand::SubmitOrderList(command) = command else {
3753                        panic!("expected SubmitOrderList command");
3754                    };
3755                    assert!(
3756                        command
3757                            .order_inits
3758                            .iter()
3759                            .all(|init| init.order_list_id == Some(command.order_list.id))
3760                    );
3761                    timeline.borrow_mut().push("command");
3762                },
3763            )
3764        };
3765        msgbus::register_trading_command_endpoint(
3766            MessagingSwitchboard::risk_engine_queue_execute(),
3767            risk_handler,
3768        );
3769
3770        let topic = format!("events.order.{}", orders[0].strategy_id());
3771        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3772
3773        strategy
3774            .submit_order_list(orders, None, None, None)
3775            .unwrap();
3776
3777        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3778
3779        let cache = strategy.cache();
3780        let cached_order1 = cache.order(&client_order_id1).unwrap();
3781        let cached_order2 = cache.order(&client_order_id2).unwrap();
3782        let order_list_id = cached_order1.order_list_id().unwrap();
3783        assert_eq!(cached_order2.order_list_id(), Some(order_list_id));
3784
3785        let event_messages = event_messages.borrow();
3786        assert_eq!(event_messages.len(), 2);
3787        let OrderEventAny::Initialized(init1) = &event_messages[0] else {
3788            panic!("expected first OrderInitialized event");
3789        };
3790        let OrderEventAny::Initialized(init2) = &event_messages[1] else {
3791            panic!("expected second OrderInitialized event");
3792        };
3793        assert_eq!(init1.order_list_id, Some(order_list_id));
3794        assert_eq!(init2.order_list_id, Some(order_list_id));
3795        assert_eq!(timeline.borrow().as_slice(), &["init1", "init2", "command"]);
3796
3797        let order_list = cache.order_list(&order_list_id).unwrap();
3798        assert_eq!(
3799            order_list.client_order_ids.as_slice(),
3800            &[client_order_id1, client_order_id2]
3801        );
3802    }
3803
3804    #[rstest]
3805    fn test_submit_order_list_routes_optional_params_to_risk() {
3806        let mut strategy = create_test_strategy();
3807        register_strategy(&mut strategy);
3808
3809        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3810            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3811        msgbus::register_trading_command_endpoint(
3812            MessagingSwitchboard::risk_engine_queue_execute(),
3813            risk_handler,
3814        );
3815
3816        let no_params_orders = vec![
3817            make_initialized_market_order("O-20250208-LIST-001"),
3818            make_initialized_market_order("O-20250208-LIST-002"),
3819        ];
3820        strategy
3821            .submit_order_list(no_params_orders, None, None, None)
3822            .unwrap();
3823
3824        let mut params = Params::new();
3825        params.insert(
3826            "routing_hint".to_string(),
3827            Value::String("prefer_batch".to_string()),
3828        );
3829        let param_orders = vec![
3830            make_initialized_market_order("O-20250208-LIST-003"),
3831            make_initialized_market_order("O-20250208-LIST-004"),
3832        ];
3833        strategy
3834            .submit_order_list(param_orders, None, None, Some(params.clone()))
3835            .unwrap();
3836
3837        let risk_messages = risk_messages.get_messages();
3838        assert_eq!(risk_messages.len(), 2);
3839        let Some(TradingCommand::SubmitOrderList(no_params_command)) = risk_messages.first() else {
3840            panic!("expected SubmitOrderList command");
3841        };
3842        let Some(TradingCommand::SubmitOrderList(param_command)) = risk_messages.get(1) else {
3843            panic!("expected SubmitOrderList command");
3844        };
3845        assert!(no_params_command.params.is_none());
3846        assert_eq!(param_command.params.as_ref(), Some(&params));
3847    }
3848
3849    #[rstest]
3850    fn test_modify_order_routes_non_emulated_orders_to_risk() {
3851        let mut strategy = create_test_strategy();
3852        register_strategy(&mut strategy);
3853
3854        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3855            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3856        msgbus::register_trading_command_endpoint(
3857            MessagingSwitchboard::risk_engine_queue_execute(),
3858            risk_handler,
3859        );
3860
3861        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3862            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3863        msgbus::register_trading_command_endpoint(
3864            MessagingSwitchboard::exec_engine_queue_execute(),
3865            exec_handler,
3866        );
3867
3868        let order = OrderAny::Market(MarketOrder::new(
3869            TraderId::from("TRADER-001"),
3870            StrategyId::from("TEST-001"),
3871            InstrumentId::from("BTCUSDT.BINANCE"),
3872            ClientOrderId::from("O-20250208-0003"),
3873            OrderSide::Buy,
3874            Quantity::from(100_000),
3875            TimeInForce::Gtc,
3876            UUID4::new(),
3877            UnixNanos::default(),
3878            false,
3879            false,
3880            None,
3881            None,
3882            None,
3883            None,
3884            None,
3885            None,
3886            None,
3887            None,
3888        ));
3889        add_order_to_cache(&strategy, &order);
3890
3891        strategy
3892            .modify_order(
3893                order.client_order_id(),
3894                Some(Quantity::from(200_000)),
3895                None,
3896                None,
3897                None,
3898                None,
3899            )
3900            .unwrap();
3901
3902        let risk_messages = risk_messages.get_messages();
3903        let exec_messages = exec_messages.get_messages();
3904
3905        assert_eq!(risk_messages.len(), 1);
3906        assert!(matches!(
3907            risk_messages.first(),
3908            Some(TradingCommand::ModifyOrder(_))
3909        ));
3910        assert!(exec_messages.is_empty());
3911    }
3912
3913    #[rstest]
3914    fn test_modify_order_routes_active_local_algorithm_order_to_algorithm() {
3915        let mut strategy = create_test_strategy();
3916        register_strategy(&mut strategy);
3917
3918        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3919            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3920        msgbus::register_trading_command_endpoint(
3921            MessagingSwitchboard::risk_engine_queue_execute(),
3922            risk_handler,
3923        );
3924        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3925            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3926        msgbus::register_trading_command_endpoint(
3927            MessagingSwitchboard::exec_engine_queue_execute(),
3928            exec_handler,
3929        );
3930        let (algo_handler, algo_messages) =
3931            get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
3932        msgbus::register_any("TWAP.execute".into(), algo_handler);
3933
3934        let order = make_initialized_algorithm_order("O-20250208-ALGO-MODIFY-001");
3935        add_order_to_cache(&strategy, &order);
3936
3937        strategy
3938            .modify_order(
3939                order.client_order_id(),
3940                Some(Quantity::from(200_000)),
3941                None,
3942                None,
3943                None,
3944                None,
3945            )
3946            .unwrap();
3947
3948        let algo_messages = algo_messages.get_messages();
3949        assert_eq!(algo_messages.len(), 1);
3950        assert!(matches!(
3951            algo_messages.first(),
3952            Some(TradingCommand::ModifyOrder(command))
3953                if command.client_order_id == order.client_order_id()
3954        ));
3955        assert!(risk_messages.get_messages().is_empty());
3956        assert!(exec_messages.get_messages().is_empty());
3957    }
3958
3959    #[rstest]
3960    fn test_modify_order_routes_accepted_algorithm_order_to_risk() {
3961        let mut strategy = create_test_strategy();
3962        register_strategy(&mut strategy);
3963
3964        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3965            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3966        msgbus::register_trading_command_endpoint(
3967            MessagingSwitchboard::risk_engine_queue_execute(),
3968            risk_handler,
3969        );
3970        let (algo_handler, algo_messages) =
3971            get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
3972        msgbus::register_any("TWAP.execute".into(), algo_handler);
3973
3974        let mut order = make_initialized_algorithm_order("O-20250208-ALGO-MODIFY-002");
3975        let account_id = AccountId::from("ACC-001");
3976        order
3977            .apply(TestOrderEventStubs::submitted(&order, account_id))
3978            .unwrap();
3979        order
3980            .apply(TestOrderEventStubs::accepted(
3981                &order,
3982                account_id,
3983                VenueOrderId::from("O-20250208-ALGO-MODIFY-002"),
3984            ))
3985            .unwrap();
3986        add_order_to_cache(&strategy, &order);
3987
3988        strategy
3989            .modify_order(
3990                order.client_order_id(),
3991                Some(Quantity::from(200_000)),
3992                None,
3993                None,
3994                None,
3995                None,
3996            )
3997            .unwrap();
3998
3999        let risk_messages = risk_messages.get_messages();
4000        assert_eq!(risk_messages.len(), 1);
4001        assert!(matches!(
4002            risk_messages.first(),
4003            Some(TradingCommand::ModifyOrder(command))
4004                if command.client_order_id == order.client_order_id()
4005        ));
4006        assert!(algo_messages.get_messages().is_empty());
4007        assert_eq!(
4008            strategy
4009                .cache()
4010                .order(&order.client_order_id())
4011                .unwrap()
4012                .status(),
4013            OrderStatus::PendingUpdate
4014        );
4015    }
4016
4017    #[rstest]
4018    fn test_modify_order_routes_emulated_order_to_order_emulator() {
4019        let mut strategy = create_test_strategy();
4020        register_strategy(&mut strategy);
4021
4022        let (emulator_handler, emulator_messages): (
4023            _,
4024            TypedIntoMessageSavingHandler<TradingCommand>,
4025        ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4026        msgbus::register_trading_command_endpoint(
4027            MessagingSwitchboard::order_emulator_execute(),
4028            emulator_handler,
4029        );
4030        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4031            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4032        msgbus::register_trading_command_endpoint(
4033            MessagingSwitchboard::risk_engine_queue_execute(),
4034            risk_handler,
4035        );
4036        let mut order = OrderTestBuilder::new(OrderType::StopMarket)
4037            .trader_id(TraderId::from("TRADER-001"))
4038            .strategy_id(StrategyId::from("TEST-001"))
4039            .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
4040            .client_order_id(ClientOrderId::from("O-20250208-EMULATED-MODIFY-001"))
4041            .side(OrderSide::Buy)
4042            .trigger_price(Price::from("51000.0"))
4043            .quantity(Quantity::from(100_000))
4044            .emulation_trigger(TriggerType::BidAsk)
4045            .build();
4046        order
4047            .apply(OrderEventAny::Emulated(
4048                OrderEmulatedSpec::builder()
4049                    .trader_id(order.trader_id())
4050                    .strategy_id(order.strategy_id())
4051                    .instrument_id(order.instrument_id())
4052                    .client_order_id(order.client_order_id())
4053                    .build(),
4054            ))
4055            .unwrap();
4056        add_order_to_cache(&strategy, &order);
4057
4058        strategy
4059            .modify_order(
4060                order.client_order_id(),
4061                None,
4062                None,
4063                Some(Price::from("52000.0")),
4064                None,
4065                None,
4066            )
4067            .unwrap();
4068
4069        let emulator_messages = emulator_messages.get_messages();
4070        assert_eq!(emulator_messages.len(), 1);
4071        assert!(matches!(
4072            emulator_messages.first(),
4073            Some(TradingCommand::ModifyOrder(command))
4074                if command.client_order_id == order.client_order_id()
4075        ));
4076        assert!(risk_messages.get_messages().is_empty());
4077    }
4078
4079    #[rstest]
4080    fn test_modify_order_routes_initialized_trigger_order_to_emulator_reentrantly() {
4081        let strategy_id = StrategyId::from("REENTRANT-001");
4082        let modified_quantity = Quantity::from(200_000);
4083        let mut strategy = InitializedModifyStrategy {
4084            core: StrategyCore::new(StrategyConfig {
4085                strategy_id: Some(strategy_id),
4086                ..Default::default()
4087            }),
4088            modified_quantity,
4089        };
4090        let trader_id = TraderId::from("TRADER-001");
4091        let clock = Rc::new(RefCell::new(TestClock::new()));
4092        let cache = Rc::new(RefCell::new(Cache::default()));
4093        let portfolio = Rc::new(RefCell::new(Portfolio::new(
4094            clock.clone(),
4095            cache.clone(),
4096            None,
4097        )));
4098        strategy
4099            .core
4100            .register(trader_id, clock, cache, portfolio)
4101            .unwrap();
4102        strategy.initialize().unwrap();
4103        strategy.start().unwrap();
4104
4105        let actor_id = strategy.actor_id().inner();
4106        register_component_actor(strategy);
4107        let order_handler = TypedHandler::from(move |event: &OrderEventAny| {
4108            let mut strategy =
4109                try_get_actor_unchecked::<InitializedModifyStrategy>(&actor_id).unwrap();
4110            strategy.handle_order_event(event.clone());
4111        });
4112        let topic = format!("events.order.{strategy_id}");
4113        msgbus::subscribe_order_events(topic.clone().into(), order_handler.clone(), None);
4114
4115        let (emulator_handler, emulator_messages): (
4116            _,
4117            TypedIntoMessageSavingHandler<TradingCommand>,
4118        ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4119        msgbus::register_trading_command_endpoint(
4120            MessagingSwitchboard::order_emulator_execute(),
4121            emulator_handler,
4122        );
4123        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4124            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4125        msgbus::register_trading_command_endpoint(
4126            MessagingSwitchboard::risk_engine_queue_execute(),
4127            risk_handler,
4128        );
4129        let (algo_handler, algo_messages) =
4130            get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
4131        msgbus::register_any("TWAP.execute".into(), algo_handler);
4132
4133        let order = OrderTestBuilder::new(OrderType::StopMarket)
4134            .trader_id(trader_id)
4135            .strategy_id(strategy_id)
4136            .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
4137            .client_order_id(ClientOrderId::from("O-REENTRANT-001"))
4138            .side(OrderSide::Buy)
4139            .trigger_price(Price::from("51000.0"))
4140            .quantity(Quantity::from(100_000))
4141            .emulation_trigger(TriggerType::BidAsk)
4142            .build();
4143        let client_order_id = order.client_order_id();
4144
4145        let mut strategy = try_get_actor_unchecked::<InitializedModifyStrategy>(&actor_id).unwrap();
4146        strategy.submit_order(order, None, None, None).unwrap();
4147        drop(strategy);
4148
4149        msgbus::unsubscribe_order_events(topic.into(), &order_handler);
4150        deregister_component(&strategy_id.inner());
4151        deregister_actor(&actor_id);
4152
4153        let emulator_messages = emulator_messages.get_messages();
4154        assert_eq!(emulator_messages.len(), 2);
4155        assert!(matches!(
4156            emulator_messages.first(),
4157            Some(TradingCommand::ModifyOrder(command))
4158                if command.client_order_id == client_order_id
4159                    && command.quantity == Some(modified_quantity)
4160        ));
4161        assert!(
4162            risk_messages
4163                .get_messages()
4164                .iter()
4165                .all(|command| !matches!(command, TradingCommand::ModifyOrder(_)))
4166        );
4167        assert!(
4168            algo_messages
4169                .get_messages()
4170                .iter()
4171                .all(|command| !matches!(command, TradingCommand::ModifyOrder(_)))
4172        );
4173        assert!(matches!(
4174            emulator_messages.get(1),
4175            Some(TradingCommand::SubmitOrder(command))
4176                if command.client_order_id == client_order_id
4177        ));
4178    }
4179
4180    #[rstest]
4181    fn test_modify_order_prefers_emulator_over_algorithm_for_emulated_algorithm_order() {
4182        let mut strategy = create_test_strategy();
4183        register_strategy(&mut strategy);
4184
4185        let (emulator_handler, emulator_messages): (
4186            _,
4187            TypedIntoMessageSavingHandler<TradingCommand>,
4188        ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4189        msgbus::register_trading_command_endpoint(
4190            MessagingSwitchboard::order_emulator_execute(),
4191            emulator_handler,
4192        );
4193        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4194            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4195        msgbus::register_trading_command_endpoint(
4196            MessagingSwitchboard::risk_engine_queue_execute(),
4197            risk_handler,
4198        );
4199        let (algo_handler, algo_messages) =
4200            get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
4201        msgbus::register_any("TWAP.execute".into(), algo_handler);
4202
4203        let mut order = OrderTestBuilder::new(OrderType::StopMarket)
4204            .trader_id(TraderId::from("TRADER-001"))
4205            .strategy_id(StrategyId::from("TEST-001"))
4206            .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
4207            .client_order_id(ClientOrderId::from("O-20250208-EMULATED-ALGO-MODIFY-001"))
4208            .side(OrderSide::Buy)
4209            .trigger_price(Price::from("51000.0"))
4210            .quantity(Quantity::from(100_000))
4211            .emulation_trigger(TriggerType::BidAsk)
4212            .exec_algorithm_id(ExecAlgorithmId::from("TWAP"))
4213            .exec_spawn_id(ClientOrderId::from("O-20250208-EMULATED-ALGO-MODIFY-001"))
4214            .build();
4215        order
4216            .apply(OrderEventAny::Emulated(
4217                OrderEmulatedSpec::builder()
4218                    .trader_id(order.trader_id())
4219                    .strategy_id(order.strategy_id())
4220                    .instrument_id(order.instrument_id())
4221                    .client_order_id(order.client_order_id())
4222                    .build(),
4223            ))
4224            .unwrap();
4225
4226        // An `Emulated` order satisfies both `is_emulated` and `is_active_local`, so this
4227        // order matches the emulator branch and the algorithm branch at once.
4228        assert!(order.is_emulated());
4229        assert!(order.is_active_local());
4230        assert!(order.exec_algorithm_id().is_some());
4231
4232        add_order_to_cache(&strategy, &order);
4233
4234        strategy
4235            .modify_order(
4236                order.client_order_id(),
4237                None,
4238                None,
4239                Some(Price::from("52000.0")),
4240                None,
4241                None,
4242            )
4243            .unwrap();
4244
4245        let emulator_messages = emulator_messages.get_messages();
4246        assert_eq!(emulator_messages.len(), 1);
4247        assert!(matches!(
4248            emulator_messages.first(),
4249            Some(TradingCommand::ModifyOrder(command))
4250                if command.client_order_id == order.client_order_id()
4251        ));
4252        assert!(algo_messages.get_messages().is_empty());
4253        assert!(risk_messages.get_messages().is_empty());
4254    }
4255
4256    #[rstest]
4257    fn test_modify_order_marks_order_pending_update_locally_before_send() {
4258        let mut strategy = create_test_strategy();
4259        register_strategy(&mut strategy);
4260
4261        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4262            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4263        msgbus::register_trading_command_endpoint(
4264            MessagingSwitchboard::risk_engine_queue_execute(),
4265            risk_handler,
4266        );
4267
4268        let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
4269            get_typed_message_saving_handler(Some(Ustr::from("events.order.pending_update")));
4270        let order = make_accepted_limit_order("O-20250208-UPDATE-001");
4271        let topic = format!("events.order.{}", order.strategy_id());
4272        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4273        add_order_to_cache(&strategy, &order);
4274
4275        strategy
4276            .modify_order(
4277                order.client_order_id(),
4278                None,
4279                Some(Price::from("51000.0")),
4280                None,
4281                None,
4282                None,
4283            )
4284            .unwrap();
4285
4286        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4287
4288        let cache = strategy.cache();
4289        let cached_order = cache.order(&order.client_order_id()).unwrap();
4290        assert_eq!(cached_order.status(), OrderStatus::PendingUpdate);
4291
4292        let risk_messages = risk_messages.get_messages();
4293        assert_eq!(risk_messages.len(), 1);
4294        assert!(matches!(
4295            risk_messages.first(),
4296            Some(TradingCommand::ModifyOrder(_))
4297        ));
4298
4299        let event_messages = event_messages.get_messages();
4300        assert_eq!(event_messages.len(), 1);
4301        assert!(matches!(
4302            event_messages.first(),
4303            Some(OrderEventAny::PendingUpdate(_))
4304        ));
4305    }
4306
4307    #[rstest]
4308    fn test_modify_orders_marks_orders_pending_update_locally_before_send() {
4309        let mut strategy = create_test_strategy();
4310        register_strategy(&mut strategy);
4311
4312        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4313            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4314        msgbus::register_trading_command_endpoint(
4315            MessagingSwitchboard::risk_engine_queue_execute(),
4316            risk_handler,
4317        );
4318
4319        let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
4320            get_typed_message_saving_handler(Some(Ustr::from("events.order.batch_pending_update")));
4321        let order1 = make_accepted_limit_order("O-20250208-BATCH-UPDATE-001");
4322        let order2 = make_accepted_limit_order("O-20250208-BATCH-UPDATE-002");
4323        let topic = format!("events.order.{}", order1.strategy_id());
4324        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4325        add_order_to_cache(&strategy, &order1);
4326        add_order_to_cache(&strategy, &order2);
4327
4328        strategy
4329            .modify_orders(
4330                vec![
4331                    (
4332                        order1.client_order_id(),
4333                        None,
4334                        Some(Price::from("51000.0")),
4335                        None,
4336                    ),
4337                    (
4338                        order2.client_order_id(),
4339                        Some(Quantity::from("2.0")),
4340                        None,
4341                        None,
4342                    ),
4343                ],
4344                None,
4345                None,
4346            )
4347            .unwrap();
4348
4349        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4350
4351        let cache = strategy.cache();
4352        let cached_order1 = cache.order(&order1.client_order_id()).unwrap();
4353        let cached_order2 = cache.order(&order2.client_order_id()).unwrap();
4354        assert_eq!(cached_order1.status(), OrderStatus::PendingUpdate);
4355        assert_eq!(cached_order2.status(), OrderStatus::PendingUpdate);
4356
4357        let risk_messages = risk_messages.get_messages();
4358        assert_eq!(risk_messages.len(), 1);
4359        let Some(TradingCommand::ModifyOrders(command)) = risk_messages.first() else {
4360            panic!("expected BatchModifyOrders command");
4361        };
4362        assert_eq!(command.modifies.len(), 2);
4363        assert_eq!(
4364            command
4365                .modifies
4366                .iter()
4367                .map(|modify| modify.client_order_id)
4368                .collect::<Vec<_>>(),
4369            vec![order1.client_order_id(), order2.client_order_id()]
4370        );
4371
4372        let event_messages = event_messages.get_messages();
4373        assert_eq!(event_messages.len(), 2);
4374        assert!(
4375            event_messages
4376                .iter()
4377                .all(|event| matches!(event, OrderEventAny::PendingUpdate(_)))
4378        );
4379    }
4380
4381    #[rstest]
4382    fn test_cancel_order_marks_order_pending_cancel_locally_before_send() {
4383        let mut strategy = create_test_strategy();
4384        register_strategy(&mut strategy);
4385
4386        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4387            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4388        msgbus::register_trading_command_endpoint(
4389            MessagingSwitchboard::exec_engine_queue_execute(),
4390            exec_handler,
4391        );
4392
4393        let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
4394            get_typed_message_saving_handler(Some(Ustr::from("events.order.pending_cancel")));
4395        let order = make_accepted_market_order("O-20250208-CANCEL-001");
4396        let topic = format!("events.order.{}", order.strategy_id());
4397        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4398        add_order_to_cache(&strategy, &order);
4399
4400        strategy
4401            .cancel_order(order.client_order_id(), None, None)
4402            .unwrap();
4403
4404        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4405
4406        let cache = strategy.cache();
4407        let cached_order = cache.order(&order.client_order_id()).unwrap();
4408        assert_eq!(cached_order.status(), OrderStatus::PendingCancel);
4409        let cache = strategy.core.cache_ref();
4410        assert!(cache.is_order_pending_cancel_local(&order.client_order_id()));
4411
4412        let exec_messages = exec_messages.get_messages();
4413        assert_eq!(exec_messages.len(), 1);
4414        assert!(matches!(
4415            exec_messages.first(),
4416            Some(TradingCommand::CancelOrder(_))
4417        ));
4418
4419        let event_messages = event_messages.get_messages();
4420        assert_eq!(event_messages.len(), 1);
4421        assert!(matches!(
4422            event_messages.first(),
4423            Some(OrderEventAny::PendingCancel(_))
4424        ));
4425    }
4426
4427    #[rstest]
4428    fn test_cancel_all_orders_strategy_only_sends_only_caller_strategy_cancels() {
4429        let mut strategy = create_test_strategy();
4430        register_strategy(&mut strategy);
4431
4432        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4433            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4434        msgbus::register_trading_command_endpoint(
4435            MessagingSwitchboard::exec_engine_queue_execute(),
4436            exec_handler,
4437        );
4438
4439        let order = make_accepted_market_order("O-20250208-CANCEL-ALL-001");
4440        let mut sibling_order = OrderTestBuilder::new(OrderType::Market)
4441            .trader_id(TraderId::from("TRADER-001"))
4442            .strategy_id(StrategyId::from("SIBLING-001"))
4443            .instrument_id(order.instrument_id())
4444            .client_order_id(ClientOrderId::from("O-20250208-CANCEL-ALL-002"))
4445            .side(OrderSide::Buy)
4446            .quantity(Quantity::from(100_000))
4447            .build();
4448        let account_id = AccountId::from("ACC-001");
4449        sibling_order
4450            .apply(TestOrderEventStubs::submitted(&sibling_order, account_id))
4451            .unwrap();
4452        sibling_order
4453            .apply(TestOrderEventStubs::accepted(
4454                &sibling_order,
4455                account_id,
4456                VenueOrderId::from("2"),
4457            ))
4458            .unwrap();
4459        add_order_to_cache(&strategy, &order);
4460        add_order_to_cache(&strategy, &sibling_order);
4461        strategy.core.cache_rc().borrow_mut().build_index();
4462
4463        strategy
4464            .cancel_all_orders(order.instrument_id(), None, None, true, None)
4465            .unwrap();
4466
4467        let messages = exec_messages.get_messages();
4468        let cache = strategy.cache();
4469        let cached_order = cache.order(&order.client_order_id()).unwrap();
4470        let cached_sibling = cache.order(&sibling_order.client_order_id()).unwrap();
4471        assert_eq!(messages.len(), 1);
4472        assert!(matches!(
4473            messages.first(),
4474            Some(TradingCommand::CancelOrder(command))
4475                if command.client_order_id == order.client_order_id()
4476        ));
4477        assert_eq!(cached_order.status(), OrderStatus::PendingCancel);
4478        assert_eq!(cached_sibling.status(), OrderStatus::Accepted);
4479    }
4480
4481    #[rstest]
4482    fn test_cancel_all_orders_strategy_only_deduplicates_emulated_inflight_order() {
4483        let mut strategy = create_test_strategy();
4484        register_strategy(&mut strategy);
4485
4486        let (emulator_handler, emulator_messages): (
4487            _,
4488            TypedIntoMessageSavingHandler<TradingCommand>,
4489        ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4490        msgbus::register_trading_command_endpoint(
4491            MessagingSwitchboard::order_emulator_execute(),
4492            emulator_handler,
4493        );
4494        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4495            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4496        msgbus::register_trading_command_endpoint(
4497            MessagingSwitchboard::exec_engine_queue_execute(),
4498            exec_handler,
4499        );
4500
4501        let client_order_id = ClientOrderId::from("O-20250208-CANCEL-ALL-EMULATED-001");
4502        let order = OrderTestBuilder::new(OrderType::StopMarket)
4503            .trader_id(TraderId::from("TRADER-001"))
4504            .strategy_id(StrategyId::from("TEST-001"))
4505            .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
4506            .client_order_id(client_order_id)
4507            .side(OrderSide::Buy)
4508            .trigger_price(Price::from("51000.0"))
4509            .quantity(Quantity::from(100_000))
4510            .emulation_trigger(TriggerType::BidAsk)
4511            .exec_algorithm_id(ExecAlgorithmId::from("ALGO-001"))
4512            .exec_spawn_id(client_order_id)
4513            .build();
4514        add_order_to_cache(&strategy, &order);
4515        strategy.core.cache_rc().borrow_mut().build_index();
4516
4517        strategy
4518            .cancel_all_orders(order.instrument_id(), None, None, true, None)
4519            .unwrap();
4520
4521        let emulator_messages = emulator_messages.get_messages();
4522        assert_eq!(emulator_messages.len(), 1);
4523        assert!(matches!(
4524            emulator_messages.first(),
4525            Some(TradingCommand::CancelOrder(command))
4526                if command.client_order_id == order.client_order_id()
4527        ));
4528        assert!(exec_messages.get_messages().is_empty());
4529    }
4530
4531    #[rstest]
4532    fn test_cancel_all_orders_without_strategy_only_sends_cancel_all_command() {
4533        let mut strategy = create_test_strategy();
4534        register_strategy(&mut strategy);
4535
4536        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4537            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4538        msgbus::register_trading_command_endpoint(
4539            MessagingSwitchboard::exec_engine_queue_execute(),
4540            exec_handler,
4541        );
4542
4543        let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
4544
4545        strategy
4546            .cancel_all_orders(instrument_id, None, None, false, None)
4547            .unwrap();
4548
4549        let messages = exec_messages.get_messages();
4550        assert_eq!(messages.len(), 1);
4551        let Some(TradingCommand::CancelAllOrders(command)) = messages.first() else {
4552            panic!("Expected a CancelAllOrders command, was {messages:?}");
4553        };
4554        assert_eq!(command.strategy_id, StrategyId::from("TEST-001"));
4555        assert_eq!(command.instrument_id, instrument_id);
4556        assert_eq!(command.client_id, None);
4557        assert_eq!(command.order_side, None);
4558        assert_eq!(command.correlation_id, Some(command.command_id));
4559        assert_eq!(command.causation_id, None);
4560    }
4561
4562    #[rstest]
4563    fn test_cancel_all_orders_without_strategy_only_delegates_local_routing_with_params() {
4564        let mut strategy = create_test_strategy();
4565        register_strategy(&mut strategy);
4566
4567        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4568            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4569        msgbus::register_trading_command_endpoint(
4570            MessagingSwitchboard::exec_engine_queue_execute(),
4571            exec_handler,
4572        );
4573        let (emulator_handler, emulator_messages): (
4574            _,
4575            TypedIntoMessageSavingHandler<TradingCommand>,
4576        ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4577        msgbus::register_trading_command_endpoint(
4578            MessagingSwitchboard::order_emulator_execute(),
4579            emulator_handler,
4580        );
4581        let exec_algorithm_id = ExecAlgorithmId::from("ALGO-001");
4582        let algorithm_endpoint = format!("{exec_algorithm_id}.execute");
4583        let (algorithm_handler, algorithm_messages) =
4584            get_any_saving_handler::<TradingCommand>(Some(Ustr::from("ALGO-001.execute")));
4585        msgbus::register_any(algorithm_endpoint.into(), algorithm_handler);
4586
4587        let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
4588        let sibling_strategy_id = StrategyId::from("SIBLING-001");
4589        let selected_client = ClientId::from("CLIENT-001");
4590        let account_id = AccountId::from("ACC-001");
4591
4592        let mut open_order = OrderTestBuilder::new(OrderType::Limit)
4593            .strategy_id(sibling_strategy_id)
4594            .instrument_id(instrument_id)
4595            .client_order_id(ClientOrderId::from("O-BROAD-OPEN-001"))
4596            .side(OrderSide::Buy)
4597            .price(Price::from("49000.0"))
4598            .quantity(Quantity::from(100_000))
4599            .build();
4600        open_order
4601            .apply(TestOrderEventStubs::submitted(&open_order, account_id))
4602            .unwrap();
4603        open_order
4604            .apply(TestOrderEventStubs::accepted(
4605                &open_order,
4606                account_id,
4607                VenueOrderId::from("V-BROAD-OPEN-001"),
4608            ))
4609            .unwrap();
4610        let emulated_order = OrderTestBuilder::new(OrderType::StopMarket)
4611            .strategy_id(sibling_strategy_id)
4612            .instrument_id(instrument_id)
4613            .client_order_id(ClientOrderId::from("O-BROAD-EMULATED-001"))
4614            .side(OrderSide::Buy)
4615            .trigger_price(Price::from("51000.0"))
4616            .quantity(Quantity::from(100_000))
4617            .emulation_trigger(TriggerType::BidAsk)
4618            .build();
4619        let algorithm_order_id = ClientOrderId::from("O-BROAD-ALGO-001");
4620        let algorithm_order = OrderTestBuilder::new(OrderType::Market)
4621            .strategy_id(sibling_strategy_id)
4622            .instrument_id(instrument_id)
4623            .client_order_id(algorithm_order_id)
4624            .side(OrderSide::Buy)
4625            .quantity(Quantity::from(100_000))
4626            .exec_algorithm_id(exec_algorithm_id)
4627            .exec_spawn_id(algorithm_order_id)
4628            .build();
4629
4630        {
4631            let cache_rc = strategy.core.cache_rc();
4632            let mut cache = cache_rc.borrow_mut();
4633            for order in [&open_order, &emulated_order, &algorithm_order] {
4634                cache
4635                    .add_order(order.clone(), None, Some(selected_client), true)
4636                    .unwrap();
4637            }
4638            cache.build_index();
4639        }
4640
4641        let mut params = Params::new();
4642        params.insert(
4643            "routing_hint".to_string(),
4644            Value::String("broad_cancel".to_string()),
4645        );
4646        strategy
4647            .cancel_all_orders(
4648                instrument_id,
4649                Some(OrderSide::Buy),
4650                Some(selected_client),
4651                false,
4652                Some(params.clone()),
4653            )
4654            .unwrap();
4655
4656        let exec_messages = exec_messages.get_messages();
4657        let Some(TradingCommand::CancelAllOrders(exec_command)) = exec_messages.first() else {
4658            panic!("Expected an execution CancelAllOrders command, was {exec_messages:?}");
4659        };
4660
4661        assert_eq!(exec_messages.len(), 1);
4662        assert!(emulator_messages.get_messages().is_empty());
4663        assert!(algorithm_messages.get_messages().is_empty());
4664        assert_eq!(exec_command.client_id, Some(selected_client));
4665        assert_eq!(exec_command.strategy_id, StrategyId::from("TEST-001"));
4666        assert_eq!(exec_command.instrument_id, instrument_id);
4667        assert_eq!(exec_command.order_side, Some(OrderSide::Buy));
4668        assert_eq!(exec_command.params.as_ref(), Some(&params));
4669        assert_eq!(exec_command.correlation_id, Some(exec_command.command_id));
4670        assert_eq!(exec_command.causation_id, None);
4671    }
4672
4673    #[rstest]
4674    fn test_cancel_all_orders_strategy_only_filters_by_client() {
4675        let mut strategy = create_test_strategy();
4676        register_strategy(&mut strategy);
4677
4678        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4679            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4680        msgbus::register_trading_command_endpoint(
4681            MessagingSwitchboard::exec_engine_queue_execute(),
4682            exec_handler,
4683        );
4684
4685        let selected_client = ClientId::from("CLIENT-001");
4686        let other_client = ClientId::from("CLIENT-002");
4687        let selected_order = make_accepted_market_order("O-20250208-CANCEL-CLIENT-001");
4688        let other_order = make_accepted_market_order("O-20250208-CANCEL-CLIENT-002");
4689        let cache_rc = strategy.core.cache_rc();
4690        {
4691            let mut cache = cache_rc.borrow_mut();
4692            cache
4693                .add_order(selected_order.clone(), None, Some(selected_client), true)
4694                .unwrap();
4695            cache
4696                .add_order(other_order.clone(), None, Some(other_client), true)
4697                .unwrap();
4698            cache.build_index();
4699        }
4700
4701        strategy
4702            .cancel_all_orders(
4703                selected_order.instrument_id(),
4704                None,
4705                Some(selected_client),
4706                true,
4707                None,
4708            )
4709            .unwrap();
4710
4711        let messages = exec_messages.get_messages();
4712        let cache = strategy.cache();
4713        let cached_selected = cache.order(&selected_order.client_order_id()).unwrap();
4714        let cached_other = cache.order(&other_order.client_order_id()).unwrap();
4715        assert_eq!(messages.len(), 1);
4716        assert!(matches!(
4717            messages.first(),
4718            Some(TradingCommand::CancelOrder(command))
4719                if command.client_id == Some(selected_client)
4720                    && command.client_order_id == selected_order.client_order_id()
4721        ));
4722        assert_eq!(cached_selected.status(), OrderStatus::PendingCancel);
4723        assert_eq!(cached_other.status(), OrderStatus::Accepted);
4724    }
4725
4726    #[rstest]
4727    fn test_cancel_all_orders_strategy_only_filters_by_side_and_preserves_params() {
4728        let mut strategy = create_test_strategy();
4729        register_strategy(&mut strategy);
4730
4731        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4732            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4733        msgbus::register_trading_command_endpoint(
4734            MessagingSwitchboard::exec_engine_queue_execute(),
4735            exec_handler,
4736        );
4737
4738        let buy_order = make_accepted_market_order("O-20250208-CANCEL-SIDE-001");
4739        let mut sell_order = OrderTestBuilder::new(OrderType::Market)
4740            .trader_id(TraderId::from("TRADER-001"))
4741            .strategy_id(StrategyId::from("TEST-001"))
4742            .instrument_id(buy_order.instrument_id())
4743            .client_order_id(ClientOrderId::from("O-20250208-CANCEL-SIDE-002"))
4744            .side(OrderSide::Sell)
4745            .quantity(Quantity::from(100_000))
4746            .build();
4747        let account_id = AccountId::from("ACC-001");
4748        sell_order
4749            .apply(TestOrderEventStubs::submitted(&sell_order, account_id))
4750            .unwrap();
4751        sell_order
4752            .apply(TestOrderEventStubs::accepted(
4753                &sell_order,
4754                account_id,
4755                VenueOrderId::from("O-20250208-CANCEL-SIDE-002"),
4756            ))
4757            .unwrap();
4758        add_order_to_cache(&strategy, &buy_order);
4759        add_order_to_cache(&strategy, &sell_order);
4760        strategy.core.cache_rc().borrow_mut().build_index();
4761
4762        let mut params = Params::new();
4763        params.insert(
4764            "routing_hint".to_string(),
4765            Value::String("strategy_only".to_string()),
4766        );
4767        strategy
4768            .cancel_all_orders(
4769                buy_order.instrument_id(),
4770                Some(OrderSide::Buy),
4771                None,
4772                true,
4773                Some(params.clone()),
4774            )
4775            .unwrap();
4776
4777        let messages = exec_messages.get_messages();
4778        let Some(TradingCommand::CancelOrder(command)) = messages.first() else {
4779            panic!("Expected a CancelOrder command, was {messages:?}");
4780        };
4781        let cache = strategy.cache();
4782        let cached_buy = cache.order(&buy_order.client_order_id()).unwrap();
4783        let cached_sell = cache.order(&sell_order.client_order_id()).unwrap();
4784        assert_eq!(messages.len(), 1);
4785        assert_eq!(command.trader_id, TraderId::from("TRADER-001"));
4786        assert_eq!(command.client_id, None);
4787        assert_eq!(command.strategy_id, StrategyId::from("TEST-001"));
4788        assert_eq!(command.instrument_id, buy_order.instrument_id());
4789        assert_eq!(command.client_order_id, buy_order.client_order_id());
4790        assert_eq!(command.venue_order_id, buy_order.venue_order_id());
4791        assert_eq!(command.params.as_ref(), Some(&params));
4792        assert_eq!(cached_buy.status(), OrderStatus::PendingCancel);
4793        assert_eq!(cached_sell.status(), OrderStatus::Accepted);
4794    }
4795
4796    #[rstest]
4797    fn test_cancel_all_orders_strategy_only_continues_after_error_and_returns_first_error() {
4798        let mut strategy = create_test_strategy();
4799        register_strategy(&mut strategy);
4800
4801        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4802            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4803        msgbus::register_trading_command_endpoint(
4804            MessagingSwitchboard::exec_engine_queue_execute(),
4805            exec_handler,
4806        );
4807
4808        let mut failing_order = make_accepted_market_order("O-20250208-CANCEL-ERROR-001");
4809        let OrderAny::Market(order) = &mut failing_order else {
4810            panic!("Expected a MarketOrder");
4811        };
4812        order.account_id = None;
4813        let succeeding_order = make_accepted_market_order("O-20250208-CANCEL-ERROR-002");
4814        add_order_to_cache(&strategy, &failing_order);
4815        add_order_to_cache(&strategy, &succeeding_order);
4816        strategy.core.cache_rc().borrow_mut().build_index();
4817
4818        let error = strategy
4819            .cancel_all_orders(failing_order.instrument_id(), None, None, true, None)
4820            .unwrap_err()
4821            .to_string();
4822
4823        let messages = exec_messages.get_messages();
4824        let cache = strategy.cache();
4825        let cached_failing = cache.order(&failing_order.client_order_id()).unwrap();
4826        let cached_succeeding = cache.order(&succeeding_order.client_order_id()).unwrap();
4827        assert_eq!(
4828            error,
4829            "Cannot generate pending cancel event for O-20250208-CANCEL-ERROR-001: \
4830             account_id is not set"
4831        );
4832        assert_eq!(messages.len(), 1);
4833        assert!(matches!(
4834            messages.first(),
4835            Some(TradingCommand::CancelOrder(command))
4836                if command.client_order_id == succeeding_order.client_order_id()
4837        ));
4838        assert_eq!(cached_failing.status(), OrderStatus::Accepted);
4839        assert_eq!(cached_succeeding.status(), OrderStatus::PendingCancel);
4840    }
4841
4842    #[rstest]
4843    fn test_cancel_orders_marks_orders_pending_cancel_locally_before_send() {
4844        let mut strategy = create_test_strategy();
4845        register_strategy(&mut strategy);
4846
4847        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4848            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4849        msgbus::register_trading_command_endpoint(
4850            MessagingSwitchboard::exec_engine_queue_execute(),
4851            exec_handler,
4852        );
4853
4854        let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
4855            get_typed_message_saving_handler(Some(Ustr::from("events.order.batch_pending_cancel")));
4856        let order1 = make_accepted_market_order("O-20250208-CANCEL-001");
4857        let order2 = make_accepted_market_order("O-20250208-CANCEL-002");
4858        let topic = format!("events.order.{}", order1.strategy_id());
4859        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4860        add_order_to_cache(&strategy, &order1);
4861        add_order_to_cache(&strategy, &order2);
4862
4863        strategy
4864            .cancel_orders(
4865                vec![order1.client_order_id(), order2.client_order_id()],
4866                None,
4867                None,
4868            )
4869            .unwrap();
4870
4871        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4872
4873        let cache = strategy.cache();
4874        let cached_order1 = cache.order(&order1.client_order_id()).unwrap();
4875        let cached_order2 = cache.order(&order2.client_order_id()).unwrap();
4876        assert_eq!(cached_order1.status(), OrderStatus::PendingCancel);
4877        assert_eq!(cached_order2.status(), OrderStatus::PendingCancel);
4878        let cache = strategy.core.cache_ref();
4879        assert!(cache.is_order_pending_cancel_local(&order1.client_order_id()));
4880        assert!(cache.is_order_pending_cancel_local(&order2.client_order_id()));
4881
4882        let exec_messages = exec_messages.get_messages();
4883        assert_eq!(exec_messages.len(), 1);
4884        let Some(TradingCommand::CancelOrders(command)) = exec_messages.first() else {
4885            panic!("expected BatchCancelOrders command");
4886        };
4887        assert_eq!(command.cancels.len(), 2);
4888
4889        let event_messages = event_messages.get_messages();
4890        assert_eq!(event_messages.len(), 2);
4891        assert!(
4892            event_messages
4893                .iter()
4894                .all(|event| matches!(event, OrderEventAny::PendingCancel(_)))
4895        );
4896    }
4897
4898    #[rstest]
4899    fn test_cancel_order_updates_own_book_status_before_send() {
4900        let mut strategy = create_test_strategy();
4901        register_strategy(&mut strategy);
4902
4903        let (exec_handler, _exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4904            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4905        msgbus::register_trading_command_endpoint(
4906            MessagingSwitchboard::exec_engine_queue_execute(),
4907            exec_handler,
4908        );
4909
4910        let order = make_accepted_limit_order("O-20250208-CANCEL-OWN-BOOK-001");
4911        add_order_to_cache_and_own_book(&strategy, &order);
4912
4913        strategy
4914            .cancel_order(order.client_order_id(), None, None)
4915            .unwrap();
4916
4917        let mut accepted = AHashSet::new();
4918        accepted.insert(OrderStatus::Accepted);
4919        let mut pending_cancel = AHashSet::new();
4920        pending_cancel.insert(OrderStatus::PendingCancel);
4921
4922        let cache = strategy.cache();
4923        let own_book = cache.own_order_book(&order.instrument_id()).unwrap();
4924        assert!(own_book.bids_as_map(Some(&accepted), None, None).is_empty());
4925        let pending_bids = own_book.bids_as_map(Some(&pending_cancel), None, None);
4926        assert_eq!(pending_bids.values().map(Vec::len).sum::<usize>(), 1);
4927    }
4928
4929    #[rstest]
4930    fn test_cancel_order_returns_error_when_not_in_cache() {
4931        let mut strategy = create_test_strategy();
4932        register_strategy(&mut strategy);
4933
4934        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4935            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4936        msgbus::register_trading_command_endpoint(
4937            MessagingSwitchboard::exec_engine_queue_execute(),
4938            exec_handler,
4939        );
4940
4941        let missing_id = ClientOrderId::from("O-MISSING");
4942        let err = strategy
4943            .cancel_order(missing_id, None, None)
4944            .expect_err("expected cancel_order to fail when order is not in cache");
4945
4946        assert_eq!(
4947            err.to_string(),
4948            format!("Cannot cancel order: {ORDER_NOT_FOUND}: {missing_id}")
4949        );
4950        assert!(exec_messages.get_messages().is_empty());
4951    }
4952
4953    #[rstest]
4954    fn test_modify_order_returns_error_when_not_in_cache() {
4955        let mut strategy = create_test_strategy();
4956        register_strategy(&mut strategy);
4957
4958        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4959            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4960        msgbus::register_trading_command_endpoint(
4961            MessagingSwitchboard::risk_engine_queue_execute(),
4962            risk_handler,
4963        );
4964
4965        let missing_id = ClientOrderId::from("O-MISSING");
4966        let err = strategy
4967            .modify_order(missing_id, Some(Quantity::from(1)), None, None, None, None)
4968            .expect_err("expected modify_order to fail when order is not in cache");
4969
4970        assert_eq!(
4971            err.to_string(),
4972            format!("Cannot modify order: {ORDER_NOT_FOUND}: {missing_id}")
4973        );
4974        assert!(risk_messages.get_messages().is_empty());
4975    }
4976
4977    #[rstest]
4978    fn test_modify_orders_returns_error_when_any_id_missing() {
4979        let mut strategy = create_test_strategy();
4980        register_strategy(&mut strategy);
4981
4982        let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4983            get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4984        msgbus::register_trading_command_endpoint(
4985            MessagingSwitchboard::risk_engine_queue_execute(),
4986            risk_handler,
4987        );
4988
4989        let order = make_accepted_limit_order("O-PRESENT");
4990        add_order_to_cache(&strategy, &order);
4991
4992        let missing_id = ClientOrderId::from("O-MISSING");
4993        let err = strategy
4994            .modify_orders(
4995                vec![
4996                    (
4997                        order.client_order_id(),
4998                        None,
4999                        Some(Price::from("51000.0")),
5000                        None,
5001                    ),
5002                    (missing_id, Some(Quantity::from("2.0")), None, None),
5003                ],
5004                None,
5005                None,
5006            )
5007            .expect_err("expected modify_orders to fail when any id is missing");
5008
5009        assert_eq!(
5010            err.to_string(),
5011            format!("Cannot modify order: {ORDER_NOT_FOUND}: {missing_id}")
5012        );
5013        assert!(risk_messages.get_messages().is_empty());
5014    }
5015
5016    #[rstest]
5017    fn test_cancel_orders_returns_error_when_any_id_missing() {
5018        let mut strategy = create_test_strategy();
5019        register_strategy(&mut strategy);
5020
5021        let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5022            get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5023        msgbus::register_trading_command_endpoint(
5024            MessagingSwitchboard::exec_engine_queue_execute(),
5025            exec_handler,
5026        );
5027
5028        let order = make_accepted_limit_order("O-PRESENT");
5029        add_order_to_cache(&strategy, &order);
5030
5031        let missing_id = ClientOrderId::from("O-MISSING");
5032        let err = strategy
5033            .cancel_orders(vec![order.client_order_id(), missing_id], None, None)
5034            .expect_err("expected cancel_orders to fail when any id is missing");
5035
5036        assert_eq!(
5037            err.to_string(),
5038            format!("Cannot cancel order: {ORDER_NOT_FOUND}: {missing_id}")
5039        );
5040        assert!(exec_messages.get_messages().is_empty());
5041    }
5042
5043    // -- GTD EXPIRY TESTS ----------------------------------------------------------------------------
5044
5045    #[rstest]
5046    fn test_has_gtd_expiry_timer_when_timer_not_set() {
5047        let mut strategy = create_test_strategy();
5048        let client_order_id = ClientOrderId::from("O-001");
5049
5050        assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5051    }
5052
5053    #[rstest]
5054    fn test_has_gtd_expiry_timer_when_timer_set() {
5055        let mut strategy = create_test_strategy();
5056        let client_order_id = ClientOrderId::from("O-001");
5057
5058        strategy
5059            .core
5060            .gtd_timers
5061            .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5062
5063        assert!(strategy.has_gtd_expiry_timer(&client_order_id));
5064    }
5065
5066    #[rstest]
5067    fn test_cancel_gtd_expiry_removes_timer() {
5068        let mut strategy = create_test_strategy();
5069        register_strategy(&mut strategy);
5070
5071        let client_order_id = ClientOrderId::from("O-001");
5072        strategy
5073            .core
5074            .gtd_timers
5075            .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5076
5077        strategy.cancel_gtd_expiry(&client_order_id);
5078
5079        assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5080    }
5081
5082    #[rstest]
5083    fn test_cancel_gtd_expiry_when_timer_not_set() {
5084        let mut strategy = create_test_strategy();
5085        register_strategy(&mut strategy);
5086
5087        let client_order_id = ClientOrderId::from("O-001");
5088
5089        strategy.cancel_gtd_expiry(&client_order_id);
5090
5091        assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5092    }
5093
5094    #[rstest]
5095    #[case::matching_order(Ustr::from("GTD-EXPIRY:O-PRESENT"))]
5096    #[case::empty_order_id(Ustr::from("GTD-EXPIRY:"))]
5097    fn test_route_time_event_ignores_unregistered_gtd_timer(#[case] timer_name: Ustr) {
5098        let mut strategy = create_test_strategy();
5099        register_strategy(&mut strategy);
5100
5101        let order = make_accepted_limit_order("O-PRESENT");
5102        let client_order_id = order.client_order_id();
5103        add_order_to_cache(&strategy, &order);
5104        let event = TimeEvent::new(
5105            timer_name,
5106            UUID4::new(),
5107            UnixNanos::default(),
5108            UnixNanos::default(),
5109        );
5110
5111        route_time_event(&mut strategy, &event);
5112
5113        let cache = strategy.core.cache_ref();
5114        let cached_order = cache.order(&client_order_id).unwrap();
5115        assert_eq!(cached_order.status(), OrderStatus::Accepted);
5116        assert!(!strategy.core.gtd_timers.contains_key(&client_order_id));
5117    }
5118
5119    #[rstest]
5120    #[case::filled(make_filled)]
5121    #[case::canceled(make_canceled)]
5122    #[case::rejected(make_rejected)]
5123    #[case::expired(make_expired)]
5124    #[case::fill_voided(make_terminal_fill_voided)]
5125    fn test_handle_order_event_cancels_gtd_timer_for_terminal_event(
5126        #[case] make_event: fn(ClientOrderId) -> OrderEventAny,
5127    ) {
5128        let mut strategy = create_test_strategy();
5129        register_strategy(&mut strategy);
5130        start_strategy(&mut strategy);
5131
5132        let client_order_id = ClientOrderId::from("O-001");
5133        strategy
5134            .core
5135            .gtd_timers
5136            .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5137
5138        strategy.handle_order_event(make_event(client_order_id));
5139
5140        assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5141    }
5142
5143    #[rstest]
5144    #[case::partial_fill(make_filled)]
5145    #[case::non_reopened_fill_void(make_terminal_fill_voided)]
5146    fn test_handle_order_event_keeps_gtd_timer_when_cached_order_remains_open(
5147        #[case] make_event: fn(ClientOrderId) -> OrderEventAny,
5148    ) {
5149        let mut strategy = create_test_strategy();
5150        register_strategy(&mut strategy);
5151        start_strategy(&mut strategy);
5152
5153        let client_order_id = ClientOrderId::from("O-001");
5154        let order = make_accepted_limit_order(client_order_id.as_str());
5155        add_order_to_cache(&strategy, &order);
5156        strategy
5157            .core
5158            .gtd_timers
5159            .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5160
5161        strategy.handle_order_event(make_event(client_order_id));
5162
5163        assert!(strategy.has_gtd_expiry_timer(&client_order_id));
5164    }
5165
5166    #[rstest]
5167    #[case::filled(make_filled)]
5168    #[case::canceled(make_canceled)]
5169    #[case::rejected(make_rejected)]
5170    #[case::expired(make_expired)]
5171    #[case::fill_voided(make_terminal_fill_voided)]
5172    fn test_handle_order_event_cancels_gtd_timer_when_stopped(
5173        #[case] make_event: fn(ClientOrderId) -> OrderEventAny,
5174    ) {
5175        let mut strategy = create_test_strategy();
5176        register_strategy(&mut strategy);
5177        start_strategy(&mut strategy);
5178
5179        let client_order_id = ClientOrderId::from("O-001");
5180        strategy
5181            .core
5182            .gtd_timers
5183            .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5184
5185        stop_strategy(&mut strategy);
5186        assert_eq!(strategy.state(), ComponentState::Stopped);
5187
5188        strategy.handle_order_event(make_event(client_order_id));
5189
5190        assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5191    }
5192
5193    #[rstest]
5194    fn test_handle_order_event_skips_gtd_cancel_for_non_terminal() {
5195        let mut strategy = create_test_strategy();
5196        register_strategy(&mut strategy);
5197        start_strategy(&mut strategy);
5198
5199        let client_order_id = ClientOrderId::from("O-001");
5200        strategy
5201            .core
5202            .gtd_timers
5203            .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5204
5205        strategy.handle_order_event(make_accepted(client_order_id));
5206
5207        assert!(strategy.has_gtd_expiry_timer(&client_order_id));
5208    }
5209
5210    #[rstest]
5211    fn test_handle_reopened_fill_void_keeps_gtd_timer() {
5212        let mut strategy = create_test_strategy();
5213        register_strategy(&mut strategy);
5214        start_strategy(&mut strategy);
5215
5216        let client_order_id = ClientOrderId::from("O-001");
5217        strategy
5218            .core
5219            .gtd_timers
5220            .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5221
5222        strategy.handle_order_event(make_fill_voided(client_order_id, true));
5223
5224        assert!(strategy.has_gtd_expiry_timer(&client_order_id));
5225    }
5226
5227    #[rstest]
5228    fn test_handle_order_event_skips_dispatch_when_stopped() {
5229        let mut strategy = create_test_strategy();
5230        register_strategy(&mut strategy);
5231        start_strategy(&mut strategy);
5232        stop_strategy(&mut strategy);
5233        assert_eq!(strategy.state(), ComponentState::Stopped);
5234
5235        strategy.handle_order_event(make_rejected(ClientOrderId::from("O-001")));
5236
5237        assert!(!strategy.on_order_event_called);
5238        assert!(!strategy.on_order_rejected_called);
5239    }
5240
5241    #[rstest]
5242    fn test_on_start_calls_reactivate_gtd_timers_when_enabled() {
5243        let config = StrategyConfig {
5244            strategy_id: Some(StrategyId::from("TEST-001")),
5245            order_id_tag: Some("001".to_string()),
5246            manage_gtd_expiry: true,
5247            ..Default::default()
5248        };
5249        let mut strategy = TestStrategy::new(config);
5250        register_strategy(&mut strategy);
5251
5252        let result = Strategy::on_start(&mut strategy);
5253        assert!(result.is_ok());
5254    }
5255
5256    #[rstest]
5257    fn test_on_start_does_not_panic_when_gtd_disabled() {
5258        let config = StrategyConfig {
5259            strategy_id: Some(StrategyId::from("TEST-001")),
5260            order_id_tag: Some("001".to_string()),
5261            manage_gtd_expiry: false,
5262            ..Default::default()
5263        };
5264        let mut strategy = TestStrategy::new(config);
5265        register_strategy(&mut strategy);
5266
5267        let result = Strategy::on_start(&mut strategy);
5268        assert!(result.is_ok());
5269    }
5270
5271    #[rstest]
5272    fn test_on_start_errors_when_strategy_id_is_not_set() {
5273        let mut strategy = TestStrategy::new(StrategyConfig::default());
5274
5275        let err = Strategy::on_start(&mut strategy).unwrap_err().to_string();
5276
5277        assert_eq!(err, "Strategy not registered: strategy_id is not set");
5278    }
5279
5280    // -- QUERY TESTS ---------------------------------------------------------------------------------
5281
5282    #[rstest]
5283    fn test_query_account_when_registered() {
5284        let mut strategy = create_test_strategy();
5285        register_strategy(&mut strategy);
5286
5287        let account_id = AccountId::from("ACC-001");
5288
5289        let result = strategy.query_account(account_id, None, None);
5290
5291        assert!(result.is_ok());
5292    }
5293
5294    #[rstest]
5295    fn test_query_account_with_client_id() {
5296        let mut strategy = create_test_strategy();
5297        register_strategy(&mut strategy);
5298
5299        let account_id = AccountId::from("ACC-001");
5300        let client_id = ClientId::from("BINANCE");
5301
5302        let result = strategy.query_account(account_id, Some(client_id), None);
5303
5304        assert!(result.is_ok());
5305    }
5306
5307    #[rstest]
5308    fn test_query_order_when_registered() {
5309        let mut strategy = create_test_strategy();
5310        register_strategy(&mut strategy);
5311
5312        let order = OrderAny::Market(MarketOrder::test_default());
5313
5314        let result = strategy.query_order(&order, None, None);
5315
5316        assert!(result.is_ok());
5317    }
5318
5319    #[rstest]
5320    fn test_query_order_with_client_id() {
5321        let mut strategy = create_test_strategy();
5322        register_strategy(&mut strategy);
5323
5324        let order = OrderAny::Market(MarketOrder::test_default());
5325        let client_id = ClientId::from("BINANCE");
5326
5327        let result = strategy.query_order(&order, Some(client_id), None);
5328
5329        assert!(result.is_ok());
5330    }
5331
5332    #[rstest]
5333    fn test_is_exiting_returns_false_by_default() {
5334        let strategy = create_test_strategy();
5335        assert!(!strategy.is_exiting());
5336    }
5337
5338    #[rstest]
5339    fn test_is_exiting_returns_true_when_set_manually() {
5340        let mut strategy = create_test_strategy();
5341        register_strategy(&mut strategy);
5342
5343        // Manually set the exiting state (as market_exit would do)
5344        strategy.core.is_exiting = true;
5345
5346        assert!(strategy.is_exiting());
5347    }
5348
5349    #[rstest]
5350    fn test_market_exit_sets_is_exiting_flag() {
5351        // Test the state changes that market_exit would make
5352        let mut strategy = create_test_strategy();
5353        register_strategy(&mut strategy);
5354
5355        assert!(!strategy.core.is_exiting);
5356
5357        // Simulate what market_exit does to the state
5358        strategy.core.is_exiting = true;
5359        strategy.core.market_exit_attempts = 0;
5360
5361        assert!(strategy.core.is_exiting);
5362        assert_eq!(strategy.core.market_exit_attempts, 0);
5363    }
5364
5365    #[rstest]
5366    fn test_market_exit_uses_config_time_in_force_and_reduce_only() {
5367        let config = StrategyConfig {
5368            strategy_id: Some(StrategyId::from("TEST-001")),
5369            order_id_tag: Some("001".to_string()),
5370            market_exit_time_in_force: TimeInForce::Ioc,
5371            market_exit_reduce_only: false,
5372            ..Default::default()
5373        };
5374        let strategy = TestStrategy::new(config);
5375
5376        assert_eq!(
5377            strategy.core.config.market_exit_time_in_force,
5378            TimeInForce::Ioc
5379        );
5380        assert!(!strategy.core.config.market_exit_reduce_only);
5381    }
5382
5383    #[rstest]
5384    fn test_market_exit_resets_attempt_counter() {
5385        let mut strategy = create_test_strategy();
5386        register_strategy(&mut strategy);
5387
5388        // Manually set attempts to simulate prior exit
5389        strategy.core.market_exit_attempts = 50;
5390
5391        // Reset via the reset method
5392        strategy.core.reset_market_exit_state();
5393
5394        assert_eq!(strategy.core.market_exit_attempts, 0);
5395    }
5396
5397    #[rstest]
5398    fn test_market_exit_second_call_returns_early_when_exiting() {
5399        let mut strategy = create_test_strategy();
5400        register_strategy(&mut strategy);
5401
5402        // First set exiting to true to simulate an in-progress exit
5403        strategy.core.is_exiting = true;
5404
5405        // Second call should return Ok and not change state
5406        let result = strategy.market_exit();
5407        assert!(result.is_ok());
5408        assert!(strategy.core.is_exiting);
5409    }
5410
5411    #[rstest]
5412    fn test_finalize_market_exit_resets_state() {
5413        let mut strategy = create_test_strategy();
5414        register_strategy(&mut strategy);
5415
5416        // Set up exiting state
5417        strategy.core.is_exiting = true;
5418        strategy.core.pending_stop = true;
5419        strategy.core.market_exit_attempts = 50;
5420
5421        strategy.finalize_market_exit();
5422
5423        assert!(!strategy.core.is_exiting);
5424        assert!(!strategy.core.pending_stop);
5425        assert_eq!(strategy.core.market_exit_attempts, 0);
5426    }
5427
5428    #[rstest]
5429    fn test_market_exit_config_defaults() {
5430        let config = StrategyConfig::default();
5431
5432        assert!(!config.manage_stop);
5433        assert_eq!(config.market_exit_interval_ms, 100);
5434        assert_eq!(config.market_exit_max_attempts, 100);
5435    }
5436
5437    #[rstest]
5438    fn test_market_exit_with_custom_config() {
5439        let config = StrategyConfig {
5440            strategy_id: Some(StrategyId::from("TEST-001")),
5441            manage_stop: true,
5442            market_exit_interval_ms: 50,
5443            market_exit_max_attempts: 200,
5444            ..Default::default()
5445        };
5446        let strategy = TestStrategy::new(config);
5447
5448        assert!(strategy.core.config.manage_stop);
5449        assert_eq!(strategy.core.config.market_exit_interval_ms, 50);
5450        assert_eq!(strategy.core.config.market_exit_max_attempts, 200);
5451    }
5452
5453    #[derive(Debug)]
5454    struct MarketExitHookTrackingStrategy {
5455        core: StrategyCore,
5456        on_market_exit_called: bool,
5457        post_market_exit_called: bool,
5458    }
5459
5460    impl MarketExitHookTrackingStrategy {
5461        fn new(config: StrategyConfig) -> Self {
5462            Self {
5463                core: StrategyCore::new(config),
5464                on_market_exit_called: false,
5465                post_market_exit_called: false,
5466            }
5467        }
5468    }
5469
5470    impl DataActor for MarketExitHookTrackingStrategy {}
5471
5472    nautilus_strategy!(MarketExitHookTrackingStrategy, {
5473        fn on_market_exit(&mut self) {
5474            self.on_market_exit_called = true;
5475        }
5476
5477        fn post_market_exit(&mut self) {
5478            self.post_market_exit_called = true;
5479        }
5480    });
5481
5482    #[rstest]
5483    fn test_market_exit_calls_on_market_exit_hook() {
5484        let config = StrategyConfig {
5485            strategy_id: Some(StrategyId::from("TEST-001")),
5486            order_id_tag: Some("001".to_string()),
5487            ..Default::default()
5488        };
5489        let mut strategy = MarketExitHookTrackingStrategy::new(config);
5490
5491        let trader_id = TraderId::from("TRADER-001");
5492        let clock = Rc::new(RefCell::new(TestClock::new()));
5493        let cache = Rc::new(RefCell::new(Cache::default()));
5494        let portfolio = Rc::new(RefCell::new(Portfolio::new(
5495            clock.clone(),
5496            cache.clone(),
5497            None,
5498        )));
5499        strategy
5500            .core
5501            .register(trader_id, clock, cache, portfolio)
5502            .unwrap();
5503        strategy.initialize().unwrap();
5504        strategy.start().unwrap();
5505
5506        let _ = strategy.market_exit();
5507
5508        assert!(strategy.on_market_exit_called);
5509    }
5510
5511    #[rstest]
5512    fn test_finalize_market_exit_calls_post_market_exit_hook() {
5513        let config = StrategyConfig {
5514            strategy_id: Some(StrategyId::from("TEST-001")),
5515            order_id_tag: Some("001".to_string()),
5516            ..Default::default()
5517        };
5518        let mut strategy = MarketExitHookTrackingStrategy::new(config);
5519
5520        let trader_id = TraderId::from("TRADER-001");
5521        let clock = Rc::new(RefCell::new(TestClock::new()));
5522        let cache = Rc::new(RefCell::new(Cache::default()));
5523        let portfolio = Rc::new(RefCell::new(Portfolio::new(
5524            clock.clone(),
5525            cache.clone(),
5526            None,
5527        )));
5528        strategy
5529            .core
5530            .register(trader_id, clock, cache, portfolio)
5531            .unwrap();
5532
5533        strategy.core.is_exiting = true;
5534        strategy.finalize_market_exit();
5535
5536        assert!(strategy.post_market_exit_called);
5537    }
5538
5539    #[derive(Debug)]
5540    struct FailingPostExitStrategy {
5541        core: StrategyCore,
5542    }
5543
5544    impl FailingPostExitStrategy {
5545        fn new(config: StrategyConfig) -> Self {
5546            Self {
5547                core: StrategyCore::new(config),
5548            }
5549        }
5550    }
5551
5552    impl DataActor for FailingPostExitStrategy {}
5553
5554    nautilus_strategy!(FailingPostExitStrategy, {
5555        fn post_market_exit(&mut self) {
5556            panic!("Simulated error in post_market_exit");
5557        }
5558    });
5559
5560    #[rstest]
5561    fn test_finalize_market_exit_handles_hook_panic() {
5562        let config = StrategyConfig {
5563            strategy_id: Some(StrategyId::from("TEST-001")),
5564            order_id_tag: Some("001".to_string()),
5565            ..Default::default()
5566        };
5567        let mut strategy = FailingPostExitStrategy::new(config);
5568
5569        let trader_id = TraderId::from("TRADER-001");
5570        let clock = Rc::new(RefCell::new(TestClock::new()));
5571        let cache = Rc::new(RefCell::new(Cache::default()));
5572        let portfolio = Rc::new(RefCell::new(Portfolio::new(
5573            clock.clone(),
5574            cache.clone(),
5575            None,
5576        )));
5577        strategy
5578            .core
5579            .register(trader_id, clock, cache, portfolio)
5580            .unwrap();
5581
5582        strategy.core.is_exiting = true;
5583        strategy.core.pending_stop = true;
5584
5585        // This should not panic - it should catch the panic in post_market_exit
5586        strategy.finalize_market_exit();
5587
5588        // State should still be reset
5589        assert!(!strategy.core.is_exiting);
5590        assert!(!strategy.core.pending_stop);
5591    }
5592
5593    #[rstest]
5594    fn test_check_market_exit_increments_attempts_before_finalizing() {
5595        let mut strategy = create_test_strategy();
5596        register_strategy(&mut strategy);
5597
5598        strategy.core.is_exiting = true;
5599        assert_eq!(strategy.core.market_exit_attempts, 0);
5600
5601        let event = TimeEvent::new(
5602            Ustr::from("MARKET_EXIT_CHECK:TEST-001"),
5603            UUID4::new(),
5604            UnixNanos::default(),
5605            UnixNanos::default(),
5606        );
5607        strategy.check_market_exit(event);
5608
5609        // With no orders/positions, check_market_exit will finalize immediately
5610        // which resets attempts to 0. This is correct behavior.
5611        // The attempt WAS incremented to 1 during the check, then reset on finalize.
5612        assert!(!strategy.core.is_exiting);
5613        assert_eq!(strategy.core.market_exit_attempts, 0);
5614    }
5615
5616    #[rstest]
5617    fn test_route_time_event_is_idempotent_when_callback_forwards() {
5618        let mut strategy = TestStrategy::new(StrategyConfig {
5619            strategy_id: Some(StrategyId::from("TEST-001")),
5620            ..Default::default()
5621        });
5622        register_strategy(&mut strategy);
5623
5624        let order = make_accepted_limit_order("O-PRESENT");
5625        add_order_to_cache(&strategy, &order);
5626        strategy.core.cache_rc().borrow_mut().build_index();
5627        strategy.core.is_exiting = true;
5628        let event = TimeEvent::new(
5629            strategy.core.market_exit_timer_name,
5630            UUID4::new(),
5631            UnixNanos::default(),
5632            UnixNanos::default(),
5633        );
5634
5635        route_time_event(&mut strategy, &event);
5636        Strategy::on_time_event(&mut strategy, &event).unwrap();
5637
5638        assert!(strategy.core.is_exiting);
5639        assert_eq!(strategy.core.market_exit_attempts, 1);
5640        assert_eq!(
5641            strategy.core.managed_time_event_last_id,
5642            Some(event.event_id)
5643        );
5644    }
5645
5646    #[rstest]
5647    fn test_route_time_event_deduplicates_overridden_managed_handlers() {
5648        let mut strategy = TimerOverrideStrategy {
5649            core: StrategyCore::new(StrategyConfig {
5650                strategy_id: Some(StrategyId::from("TEST-001")),
5651                ..Default::default()
5652            }),
5653            gtd_expiries: 0,
5654            market_exit_checks: 0,
5655        };
5656        let market_event = TimeEvent::new(
5657            strategy.core.market_exit_timer_name,
5658            UUID4::new(),
5659            UnixNanos::default(),
5660            UnixNanos::default(),
5661        );
5662
5663        route_time_event(&mut strategy, &market_event);
5664        DataActor::on_time_event(&mut strategy, &market_event).unwrap();
5665
5666        let client_order_id = ClientOrderId::from("O-001");
5667        let gtd_timer_name = Ustr::from("GTD-EXPIRY:O-001");
5668        strategy
5669            .core
5670            .gtd_timers
5671            .insert(client_order_id, gtd_timer_name);
5672        let gtd_event = TimeEvent::new(
5673            gtd_timer_name,
5674            UUID4::new(),
5675            UnixNanos::default(),
5676            UnixNanos::default(),
5677        );
5678
5679        route_time_event(&mut strategy, &gtd_event);
5680        DataActor::on_time_event(&mut strategy, &gtd_event).unwrap();
5681
5682        assert_eq!(strategy.market_exit_checks, 1);
5683        assert_eq!(strategy.gtd_expiries, 1);
5684        assert!(strategy.core.gtd_timers.contains_key(&client_order_id));
5685        assert_eq!(
5686            strategy.core.managed_time_event_last_id,
5687            Some(gtd_event.event_id)
5688        );
5689    }
5690
5691    #[rstest]
5692    fn test_route_time_event_ignores_unowned_market_exit_timer() {
5693        let mut strategy = TestStrategy::new(StrategyConfig {
5694            strategy_id: Some(StrategyId::from("TEST-001")),
5695            ..Default::default()
5696        });
5697        register_strategy(&mut strategy);
5698
5699        let order = make_accepted_limit_order("O-PRESENT");
5700        add_order_to_cache(&strategy, &order);
5701        strategy.core.cache_rc().borrow_mut().build_index();
5702        strategy.core.is_exiting = true;
5703        let event = TimeEvent::new(
5704            Ustr::from("MARKET_EXIT_CHECK:OTHER-001"),
5705            UUID4::new(),
5706            UnixNanos::default(),
5707            UnixNanos::default(),
5708        );
5709
5710        route_time_event(&mut strategy, &event);
5711
5712        assert!(strategy.core.is_exiting);
5713        assert_eq!(strategy.core.market_exit_attempts, 0);
5714        assert_eq!(strategy.core.managed_time_event_last_id, None);
5715    }
5716
5717    #[rstest]
5718    fn test_check_market_exit_finalizes_when_max_attempts_reached() {
5719        let config = StrategyConfig {
5720            strategy_id: Some(StrategyId::from("TEST-001")),
5721            order_id_tag: Some("001".to_string()),
5722            market_exit_max_attempts: 3,
5723            ..Default::default()
5724        };
5725        let mut strategy = TestStrategy::new(config);
5726        register_strategy(&mut strategy);
5727
5728        strategy.core.is_exiting = true;
5729        strategy.core.market_exit_attempts = 2; // One below max
5730
5731        let event = TimeEvent::new(
5732            Ustr::from("MARKET_EXIT_CHECK:TEST-001"),
5733            UUID4::new(),
5734            UnixNanos::default(),
5735            UnixNanos::default(),
5736        );
5737        strategy.check_market_exit(event);
5738
5739        // Should have finalized since attempts >= max_attempts
5740        assert!(!strategy.core.is_exiting);
5741        assert_eq!(strategy.core.market_exit_attempts, 0);
5742    }
5743
5744    #[rstest]
5745    fn test_check_market_exit_finalizes_when_no_orders_or_positions() {
5746        let mut strategy = create_test_strategy();
5747        register_strategy(&mut strategy);
5748
5749        strategy.core.is_exiting = true;
5750
5751        let event = TimeEvent::new(
5752            Ustr::from("MARKET_EXIT_CHECK:TEST-001"),
5753            UUID4::new(),
5754            UnixNanos::default(),
5755            UnixNanos::default(),
5756        );
5757        strategy.check_market_exit(event);
5758
5759        // Should have finalized since there are no orders or positions
5760        assert!(!strategy.core.is_exiting);
5761    }
5762
5763    #[rstest]
5764    fn test_market_exit_timer_name_format() {
5765        let config = StrategyConfig {
5766            strategy_id: Some(StrategyId::from("MY-STRATEGY-001")),
5767            ..Default::default()
5768        };
5769        let strategy = TestStrategy::new(config);
5770
5771        assert_eq!(
5772            strategy.core.market_exit_timer_name.as_str(),
5773            "MARKET_EXIT_CHECK:MY-STRATEGY-001"
5774        );
5775    }
5776
5777    #[rstest]
5778    fn test_reset_market_exit_state() {
5779        let mut strategy = create_test_strategy();
5780
5781        strategy.core.is_exiting = true;
5782        strategy.core.pending_stop = true;
5783        strategy.core.market_exit_attempts = 50;
5784
5785        strategy.core.reset_market_exit_state();
5786
5787        assert!(!strategy.core.is_exiting);
5788        assert!(!strategy.core.pending_stop);
5789        assert_eq!(strategy.core.market_exit_attempts, 0);
5790    }
5791
5792    #[rstest]
5793    fn test_cancel_market_exit_resets_state_without_hooks() {
5794        let config = StrategyConfig {
5795            strategy_id: Some(StrategyId::from("TEST-001")),
5796            order_id_tag: Some("001".to_string()),
5797            ..Default::default()
5798        };
5799        let mut strategy = MarketExitHookTrackingStrategy::new(config);
5800
5801        let trader_id = TraderId::from("TRADER-001");
5802        let clock = Rc::new(RefCell::new(TestClock::new()));
5803        let cache = Rc::new(RefCell::new(Cache::default()));
5804        let portfolio = Rc::new(RefCell::new(Portfolio::new(
5805            clock.clone(),
5806            cache.clone(),
5807            None,
5808        )));
5809        strategy
5810            .core
5811            .register(trader_id, clock, cache, portfolio)
5812            .unwrap();
5813
5814        // Set up exiting state
5815        strategy.core.is_exiting = true;
5816        strategy.core.pending_stop = true;
5817        strategy.core.market_exit_attempts = 50;
5818
5819        // Call cancel_market_exit
5820        strategy.cancel_market_exit();
5821
5822        // State should be reset
5823        assert!(!strategy.core.is_exiting);
5824        assert!(!strategy.core.pending_stop);
5825        assert_eq!(strategy.core.market_exit_attempts, 0);
5826
5827        // Hooks should NOT have been called
5828        assert!(!strategy.on_market_exit_called);
5829        assert!(!strategy.post_market_exit_called);
5830    }
5831
5832    #[rstest]
5833    fn test_market_exit_returns_early_when_not_running() {
5834        let mut strategy = create_test_strategy();
5835        register_strategy(&mut strategy);
5836
5837        // State is not Running (default is PreInitialized)
5838        assert!(!strategy.is_running());
5839
5840        let result = strategy.market_exit();
5841
5842        // Should return Ok but not set is_exiting
5843        assert!(result.is_ok());
5844        assert!(!strategy.core.is_exiting);
5845    }
5846
5847    #[rstest]
5848    fn test_stop_with_manage_stop_false_cleans_up_active_exit() {
5849        let config = StrategyConfig {
5850            strategy_id: Some(StrategyId::from("TEST-001")),
5851            order_id_tag: Some("001".to_string()),
5852            manage_stop: false,
5853            ..Default::default()
5854        };
5855        let mut strategy = TestStrategy::new(config);
5856        register_strategy(&mut strategy);
5857
5858        // Simulate an active market exit
5859        strategy.core.is_exiting = true;
5860        strategy.core.market_exit_attempts = 5;
5861
5862        // Call stop
5863        let should_proceed = Strategy::stop(&mut strategy);
5864
5865        // Should clean up state and allow stop to proceed
5866        assert!(should_proceed);
5867        assert!(!strategy.core.is_exiting);
5868        assert_eq!(strategy.core.market_exit_attempts, 0);
5869    }
5870
5871    #[rstest]
5872    fn test_stop_with_manage_stop_true_defers_when_running() {
5873        let config = StrategyConfig {
5874            strategy_id: Some(StrategyId::from("TEST-001")),
5875            order_id_tag: Some("001".to_string()),
5876            manage_stop: true,
5877            ..Default::default()
5878        };
5879        let mut strategy = TestStrategy::new(config);
5880
5881        // Custom setup with a default callback so timer scheduling succeeds
5882        let trader_id = TraderId::from("TRADER-001");
5883        let clock = Rc::new(RefCell::new(TestClock::new()));
5884        clock
5885            .borrow_mut()
5886            .register_default_handler(TimeEventCallback::from(|_event: TimeEvent| {}));
5887        let cache = Rc::new(RefCell::new(Cache::default()));
5888        let portfolio = Rc::new(RefCell::new(Portfolio::new(
5889            clock.clone(),
5890            cache.clone(),
5891            None,
5892        )));
5893        strategy
5894            .core
5895            .register(trader_id, clock, cache, portfolio)
5896            .unwrap();
5897        strategy.initialize().unwrap();
5898        strategy.start().unwrap();
5899
5900        let should_proceed = Strategy::stop(&mut strategy);
5901
5902        // Should set pending_stop and defer
5903        assert!(!should_proceed);
5904        assert!(strategy.core.pending_stop);
5905    }
5906
5907    #[rstest]
5908    fn test_stop_with_manage_stop_true_returns_early_if_pending() {
5909        let config = StrategyConfig {
5910            strategy_id: Some(StrategyId::from("TEST-001")),
5911            order_id_tag: Some("001".to_string()),
5912            manage_stop: true,
5913            ..Default::default()
5914        };
5915        let mut strategy = TestStrategy::new(config);
5916        register_strategy(&mut strategy);
5917        start_strategy(&mut strategy);
5918        strategy.core.pending_stop = true;
5919
5920        // Call stop again
5921        let should_proceed = Strategy::stop(&mut strategy);
5922
5923        // Should return early without changing state
5924        assert!(!should_proceed);
5925        assert!(strategy.core.pending_stop);
5926    }
5927
5928    #[rstest]
5929    fn test_stop_with_manage_stop_true_proceeds_when_not_running() {
5930        let config = StrategyConfig {
5931            strategy_id: Some(StrategyId::from("TEST-001")),
5932            order_id_tag: Some("001".to_string()),
5933            manage_stop: true,
5934            ..Default::default()
5935        };
5936        let mut strategy = TestStrategy::new(config);
5937        register_strategy(&mut strategy);
5938
5939        // State is not Running (default)
5940        assert!(!strategy.is_running());
5941
5942        let should_proceed = Strategy::stop(&mut strategy);
5943
5944        // Should proceed with stop
5945        assert!(should_proceed);
5946    }
5947
5948    #[rstest]
5949    fn test_finalize_market_exit_stops_strategy_when_pending() {
5950        let config = StrategyConfig {
5951            strategy_id: Some(StrategyId::from("TEST-001")),
5952            order_id_tag: Some("001".to_string()),
5953            ..Default::default()
5954        };
5955        let mut strategy = TestStrategy::new(config);
5956        register_strategy(&mut strategy);
5957        start_strategy(&mut strategy);
5958
5959        // Simulate a market exit with pending stop
5960        strategy.core.is_exiting = true;
5961        strategy.core.pending_stop = true;
5962
5963        strategy.finalize_market_exit();
5964
5965        // Should have transitioned to Stopped
5966        assert_eq!(strategy.state(), ComponentState::Stopped);
5967        assert!(!strategy.core.is_exiting);
5968        assert!(!strategy.core.pending_stop);
5969    }
5970
5971    #[rstest]
5972    fn test_finalize_market_exit_stays_running_when_not_pending() {
5973        let config = StrategyConfig {
5974            strategy_id: Some(StrategyId::from("TEST-001")),
5975            order_id_tag: Some("001".to_string()),
5976            ..Default::default()
5977        };
5978        let mut strategy = TestStrategy::new(config);
5979        register_strategy(&mut strategy);
5980        start_strategy(&mut strategy);
5981
5982        // Simulate a market exit without pending stop
5983        strategy.core.is_exiting = true;
5984        strategy.core.pending_stop = false;
5985
5986        strategy.finalize_market_exit();
5987
5988        // Should stay Running
5989        assert_eq!(strategy.state(), ComponentState::Running);
5990        assert!(!strategy.core.is_exiting);
5991    }
5992
5993    #[rstest]
5994    fn test_submit_order_denied_during_market_exit_when_not_reduce_only() {
5995        let mut strategy = create_test_strategy();
5996        register_strategy(&mut strategy);
5997        start_strategy(&mut strategy);
5998        strategy.core.is_exiting = true;
5999
6000        let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
6001            get_typed_message_saving_handler(Some(Ustr::from("events.order.denied")));
6002        let order = OrderAny::Market(MarketOrder::new(
6003            TraderId::from("TRADER-001"),
6004            StrategyId::from("TEST-001"),
6005            InstrumentId::from("BTCUSDT.BINANCE"),
6006            ClientOrderId::from("O-20250208-0001"),
6007            OrderSide::Buy,
6008            Quantity::from(100_000),
6009            TimeInForce::Gtc,
6010            UUID4::new(),
6011            UnixNanos::default(),
6012            false, // not reduce_only
6013            false,
6014            None,
6015            None,
6016            None,
6017            None,
6018            None,
6019            None,
6020            None,
6021            None,
6022        ));
6023        let topic = format!("events.order.{}", order.strategy_id());
6024        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
6025        let client_order_id = order.client_order_id();
6026        let result = strategy.submit_order(order.clone(), None, None, None);
6027
6028        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
6029
6030        assert!(result.is_ok());
6031        let cache = strategy.cache();
6032        let cached_order = cache.order(&client_order_id).unwrap();
6033        assert_eq!(cached_order.status(), OrderStatus::Denied);
6034
6035        let event_messages = event_messages.get_messages();
6036        assert_eq!(event_messages.len(), 2);
6037        assert_eq!(
6038            event_messages[0],
6039            OrderEventAny::Initialized(order.init_event().clone())
6040        );
6041        let OrderEventAny::Denied(denied) = &event_messages[1] else {
6042            panic!("expected OrderDenied event");
6043        };
6044        assert_eq!(denied.reason, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
6045    }
6046
6047    #[rstest]
6048    fn test_submit_order_list_denied_during_market_exit_publishes_init_then_denied_events() {
6049        let mut strategy = create_test_strategy();
6050        register_strategy(&mut strategy);
6051        start_strategy(&mut strategy);
6052        strategy.core.is_exiting = true;
6053
6054        let orders = vec![
6055            make_initialized_market_order("O-20250208-LIST-DENY-001"),
6056            make_initialized_market_order("O-20250208-LIST-DENY-002"),
6057        ];
6058        let client_order_id1 = orders[0].client_order_id();
6059        let client_order_id2 = orders[1].client_order_id();
6060        let cache_rc = strategy.core.cache_rc();
6061        let timeline = Rc::new(RefCell::new(Vec::new()));
6062        let event_messages = Rc::new(RefCell::new(Vec::new()));
6063
6064        let event_handler = {
6065            let event_messages = event_messages.clone();
6066            let timeline = timeline.clone();
6067            TypedHandler::from_with_id("events.order.list_denied", move |event: &OrderEventAny| {
6068                match event {
6069                    OrderEventAny::Initialized(e) if e.client_order_id == client_order_id1 => {
6070                        assert!(cache_rc.borrow().order_exists(&client_order_id1));
6071                        timeline.borrow_mut().push("init1");
6072                    }
6073                    OrderEventAny::Initialized(e) if e.client_order_id == client_order_id2 => {
6074                        assert!(cache_rc.borrow().order_exists(&client_order_id2));
6075                        timeline.borrow_mut().push("init2");
6076                    }
6077                    OrderEventAny::Denied(e) if e.client_order_id == client_order_id1 => {
6078                        assert_eq!(e.reason, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
6079                        let cache = cache_rc.borrow();
6080                        let cached_order = cache.order(&client_order_id1).unwrap();
6081                        assert_eq!(cached_order.status(), OrderStatus::Denied);
6082                        timeline.borrow_mut().push("denied1");
6083                    }
6084                    OrderEventAny::Denied(e) if e.client_order_id == client_order_id2 => {
6085                        assert_eq!(e.reason, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
6086                        let cache = cache_rc.borrow();
6087                        let cached_order = cache.order(&client_order_id2).unwrap();
6088                        assert_eq!(cached_order.status(), OrderStatus::Denied);
6089                        timeline.borrow_mut().push("denied2");
6090                    }
6091                    _ => panic!("unexpected order event {event:?}"),
6092                }
6093                event_messages.borrow_mut().push(event.clone());
6094            })
6095        };
6096        let risk_handler = {
6097            let timeline = timeline.clone();
6098            TypedIntoHandler::from_with_id(
6099                "RiskEngine.queue_execute",
6100                move |_command: TradingCommand| {
6101                    timeline.borrow_mut().push("command");
6102                },
6103            )
6104        };
6105        msgbus::register_trading_command_endpoint(
6106            MessagingSwitchboard::risk_engine_queue_execute(),
6107            risk_handler,
6108        );
6109
6110        let topic = format!("events.order.{}", orders[0].strategy_id());
6111        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
6112        let result = strategy.submit_order_list(orders.clone(), None, None, None);
6113
6114        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
6115
6116        assert!(result.is_ok());
6117
6118        let cache = strategy.cache();
6119        let cached_order1 = cache.order(&client_order_id1).unwrap();
6120        let cached_order2 = cache.order(&client_order_id2).unwrap();
6121        assert_eq!(cached_order1.status(), OrderStatus::Denied);
6122        assert_eq!(cached_order2.status(), OrderStatus::Denied);
6123
6124        let event_messages = event_messages.borrow();
6125        assert_eq!(event_messages.len(), 4);
6126        assert_eq!(
6127            event_messages[0],
6128            OrderEventAny::Initialized(orders[0].init_event().clone())
6129        );
6130        assert!(matches!(
6131            &event_messages[1],
6132            OrderEventAny::Denied(e)
6133                if e.client_order_id == client_order_id1
6134                    && e.reason == Ustr::from("MARKET_EXIT_IN_PROGRESS")
6135        ));
6136        assert_eq!(
6137            event_messages[2],
6138            OrderEventAny::Initialized(orders[1].init_event().clone())
6139        );
6140        assert!(matches!(
6141            &event_messages[3],
6142            OrderEventAny::Denied(e)
6143                if e.client_order_id == client_order_id2
6144                    && e.reason == Ustr::from("MARKET_EXIT_IN_PROGRESS")
6145        ));
6146        assert_eq!(
6147            timeline.borrow().as_slice(),
6148            &["init1", "denied1", "init2", "denied2"]
6149        );
6150    }
6151
6152    #[rstest]
6153    fn test_submit_order_list_market_exit_rejects_non_initialized_without_events() {
6154        let mut strategy = create_test_strategy();
6155        register_strategy(&mut strategy);
6156        start_strategy(&mut strategy);
6157        strategy.core.is_exiting = true;
6158
6159        let order = make_accepted_market_order("O-20250208-LIST-DENY-ACCEPTED");
6160        let topic = format!("events.order.{}", order.strategy_id());
6161        let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
6162            get_typed_message_saving_handler(Some(Ustr::from("events.order.list_invalid")));
6163
6164        msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
6165        let result = strategy.submit_order_list(vec![order], None, None, None);
6166
6167        msgbus::unsubscribe_order_events(topic.into(), &event_handler);
6168
6169        assert!(result.is_err());
6170        assert!(
6171            result
6172                .unwrap_err()
6173                .to_string()
6174                .contains("expected INITIALIZED")
6175        );
6176        assert!(event_messages.get_messages().is_empty());
6177    }
6178
6179    #[rstest]
6180    fn test_submit_order_list_rejects_mixed_venues_with_friendly_error() {
6181        let mut strategy = create_test_strategy();
6182        register_strategy(&mut strategy);
6183        start_strategy(&mut strategy);
6184
6185        let binance_order = make_initialized_market_order("O-MIXED-VENUE-001");
6186        let bybit_order = OrderAny::Market(MarketOrder::new(
6187            TraderId::from("TRADER-001"),
6188            StrategyId::from("TEST-001"),
6189            InstrumentId::from("BTCUSDT.BYBIT"),
6190            ClientOrderId::from("O-MIXED-VENUE-002"),
6191            OrderSide::Buy,
6192            Quantity::from(100_000),
6193            TimeInForce::Gtc,
6194            UUID4::new(),
6195            UnixNanos::default(),
6196            false,
6197            false,
6198            None,
6199            None,
6200            None,
6201            None,
6202            None,
6203            None,
6204            None,
6205            None,
6206        ));
6207
6208        let result = strategy.submit_order_list(vec![binance_order, bybit_order], None, None, None);
6209
6210        let err = result.unwrap_err();
6211        let msg = err.to_string();
6212        assert!(
6213            msg.contains("OrderList denied: orders must share the same venue"),
6214            "unexpected error: {msg}",
6215        );
6216        assert!(msg.contains("BINANCE"), "expected BINANCE in error: {msg}");
6217        assert!(msg.contains("BYBIT"), "expected BYBIT in error: {msg}");
6218    }
6219
6220    #[rstest]
6221    fn test_submit_order_allowed_during_market_exit_when_reduce_only() {
6222        let mut strategy = create_test_strategy();
6223        register_strategy(&mut strategy);
6224        start_strategy(&mut strategy);
6225        strategy.core.is_exiting = true;
6226
6227        let order = OrderAny::Market(MarketOrder::new(
6228            TraderId::from("TRADER-001"),
6229            StrategyId::from("TEST-001"),
6230            InstrumentId::from("BTCUSDT.BINANCE"),
6231            ClientOrderId::from("O-20250208-0001"),
6232            OrderSide::Buy,
6233            Quantity::from(100_000),
6234            TimeInForce::Gtc,
6235            UUID4::new(),
6236            UnixNanos::default(),
6237            true, // reduce_only
6238            false,
6239            None,
6240            None,
6241            None,
6242            None,
6243            None,
6244            None,
6245            None,
6246            None,
6247        ));
6248        let client_order_id = order.client_order_id();
6249        let result = strategy.submit_order(order, None, None, None);
6250
6251        assert!(result.is_ok());
6252        let cache = strategy.cache();
6253        let cached_order = cache.order(&client_order_id).unwrap();
6254        assert_ne!(cached_order.status(), OrderStatus::Denied);
6255    }
6256
6257    #[rstest]
6258    fn test_submit_order_allowed_during_market_exit_when_tagged() {
6259        let mut strategy = create_test_strategy();
6260        register_strategy(&mut strategy);
6261        start_strategy(&mut strategy);
6262        strategy.core.is_exiting = true;
6263
6264        let order = OrderAny::Market(MarketOrder::new(
6265            TraderId::from("TRADER-001"),
6266            StrategyId::from("TEST-001"),
6267            InstrumentId::from("BTCUSDT.BINANCE"),
6268            ClientOrderId::from("O-20250208-0002"),
6269            OrderSide::Buy,
6270            Quantity::from(100_000),
6271            TimeInForce::Gtc,
6272            UUID4::new(),
6273            UnixNanos::default(),
6274            false, // not reduce_only
6275            false,
6276            None,
6277            None,
6278            None,
6279            None,
6280            None,
6281            None,
6282            None,
6283            Some(vec![Ustr::from("MARKET_EXIT")]),
6284        ));
6285        let client_order_id = order.client_order_id();
6286        let result = strategy.submit_order(order, None, None, None);
6287
6288        assert!(result.is_ok());
6289        let cache = strategy.cache();
6290        let cached_order = cache.order(&client_order_id).unwrap();
6291        assert_ne!(cached_order.status(), OrderStatus::Denied);
6292    }
6293
6294    #[derive(Debug)]
6295    struct MacroTestSimple {
6296        core: StrategyCore,
6297    }
6298
6299    nautilus_strategy!(MacroTestSimple);
6300
6301    impl DataActor for MacroTestSimple {}
6302
6303    #[derive(Debug)]
6304    struct MacroTestWithHooks {
6305        core: StrategyCore,
6306    }
6307
6308    nautilus_strategy!(MacroTestWithHooks, {
6309        fn on_order_rejected(&mut self, _event: OrderRejected) {}
6310    });
6311
6312    impl DataActor for MacroTestWithHooks {}
6313
6314    #[derive(Debug)]
6315    struct MacroTestCustomField {
6316        inner: StrategyCore,
6317    }
6318
6319    nautilus_strategy!(MacroTestCustomField, inner, {
6320        fn external_order_claims(&self) -> Option<Vec<InstrumentId>> {
6321            None
6322        }
6323    });
6324
6325    impl DataActor for MacroTestCustomField {}
6326
6327    #[rstest]
6328    fn test_strategy_behavior_does_not_require_native_core_access() {
6329        fn assert_strategy<T: Strategy + DataActor>() {}
6330
6331        assert_strategy::<CoreFreeStrategy>();
6332
6333        let mut strategy = CoreFreeStrategy { started: false };
6334        DataActor::on_start(&mut strategy).unwrap();
6335
6336        assert!(strategy.started);
6337    }
6338
6339    #[rstest]
6340    fn test_nautilus_strategy_macro_forms() {
6341        let config = StrategyConfig {
6342            strategy_id: Some(StrategyId::from("MACRO-001")),
6343            order_id_tag: Some("001".to_string()),
6344            ..Default::default()
6345        };
6346
6347        let simple = MacroTestSimple {
6348            core: StrategyCore::new(config.clone()),
6349        };
6350        assert_eq!(simple.strategy_id(), config.strategy_id);
6351        assert_eq!(simple.config().order_id_tag, config.order_id_tag);
6352        assert_eq!(simple.actor_id(), ActorId::from("MACRO-001"));
6353
6354        let hooks = MacroTestWithHooks {
6355            core: StrategyCore::new(config.clone()),
6356        };
6357        assert_eq!(hooks.strategy_id(), config.strategy_id);
6358        assert_eq!(hooks.config().order_id_tag, config.order_id_tag);
6359        assert_eq!(hooks.actor_id(), ActorId::from("MACRO-001"));
6360
6361        let custom = MacroTestCustomField {
6362            inner: StrategyCore::new(config.clone()),
6363        };
6364        assert_eq!(custom.strategy_id(), config.strategy_id);
6365        assert_eq!(custom.config().order_id_tag, config.order_id_tag);
6366        assert_eq!(custom.actor_id(), ActorId::from("MACRO-001"));
6367        assert!(custom.external_order_claims().is_none());
6368    }
6369}