Skip to main content

nautilus_trading/algorithm/
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
16//! Execution algorithm infrastructure for order slicing and execution optimization.
17//!
18//! This module provides the [`ExecutionAlgorithm`] trait and supporting infrastructure
19//! for implementing algorithms like TWAP (Time-Weighted Average Price) and VWAP
20//! (Volume-Weighted Average Price) that slice large orders into smaller child orders.
21//!
22//! # Architecture
23//!
24//! Execution algorithms extend [`DataActor`] (not [`Strategy`](super::Strategy)) because:
25//! - They don't own positions (the parent Strategy does).
26//! - Spawned orders carry the parent Strategy's ID, not the algorithm's ID.
27//! - They act as order processors/transformers, not position managers.
28//!
29//! # Order Flow
30//!
31//! 1. A Strategy submits an order with `exec_algorithm_id` set.
32//! 2. The order is routed to the algorithm's `{id}.execute` endpoint.
33//! 3. The algorithm receives the order via `on_order()`.
34//! 4. The algorithm spawns child orders using `spawn_market()`, `spawn_limit()`, etc.
35//! 5. Spawned orders are submitted through the `RiskEngine`.
36//! 6. The algorithm receives fill events and manages remaining quantity.
37
38pub mod config;
39pub mod core;
40pub mod twap;
41
42pub use core::{ExecutionAlgorithmCore, ExecutionAlgorithmNative, StrategyEventHandlers};
43
44pub use config::{ExecutionAlgorithmConfig, ImportableExecAlgorithmConfig};
45use nautilus_common::{
46    actor::{DataActor, DataActorNative, registry::try_get_actor_unchecked},
47    enums::ComponentState,
48    logging::{CMD, EVT, RECV, SEND},
49    messages::execution::{CancelOrder, ModifyOrder, SubmitOrder, TradingCommand},
50    msgbus::{self, MessagingSwitchboard, TypedHandler},
51    timer::TimeEvent,
52};
53use nautilus_core::{UUID4, UnixNanos};
54use nautilus_model::{
55    enums::{OrderStatus, TimeInForce, TriggerType},
56    events::{
57        OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied, OrderEmulated,
58        OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled, OrderInitialized,
59        OrderModifyRejected, OrderPendingCancel, OrderPendingUpdate, OrderRejected, OrderReleased,
60        OrderSubmitted, OrderTriggered, OrderUpdated, PositionChanged, PositionClosed,
61        PositionEvent, PositionOpened,
62    },
63    identifiers::{AccountId, ClientId, ExecAlgorithmId, PositionId, StrategyId, TraderId},
64    orders::{LimitOrder, MarketOrder, MarketToLimitOrder, Order, OrderAny, OrderError, OrderList},
65    types::{Price, Quantity},
66};
67pub use twap::{TwapAlgorithm, TwapAlgorithmConfig};
68use ustr::Ustr;
69
70/// Core trait for implementing execution algorithms in NautilusTrader.
71///
72/// Execution algorithms are specialized [`DataActor`]s that receive orders from strategies
73/// and execute them by spawning child orders. They are used for order slicing algorithms
74/// like TWAP and VWAP.
75///
76/// # Key Capabilities
77///
78/// - All [`DataActor`] capabilities (data subscriptions, event handling, timers)
79/// - Order spawning (market, limit, market-to-limit)
80/// - Order lifecycle management (submit, modify, cancel)
81/// - Event filtering for algorithm-owned orders
82///
83/// # Implementation
84///
85/// Use the `nautilus_execution_algorithm!` macro to generate the native runtime
86/// wiring and `ExecutionAlgorithm` implementation, including the required
87/// `on_order()` method. Normal execution algorithm logic should call facade
88/// methods such as `submit_order()`, `spawn_market()`, and
89/// `unsubscribe_all_strategy_events()`. Native runtime code that needs the
90/// internal core should use [`ExecutionAlgorithmNative`].
91pub trait ExecutionAlgorithm: DataActor {
92    /// Returns the execution algorithm ID.
93    fn id(&self) -> ExecAlgorithmId
94    where
95        Self: ExecutionAlgorithmNative,
96    {
97        ExecutionAlgorithmNative::exec_algorithm_core(self).exec_algorithm_id
98    }
99
100    /// Executes a trading command.
101    ///
102    /// This is the main entry point for commands routed to the algorithm.
103    /// Dispatches to the appropriate handler based on command type.
104    ///
105    /// Commands are only processed when the algorithm is in `Running` state.
106    ///
107    /// # Errors
108    ///
109    /// Returns an error if command handling fails.
110    fn execute(&mut self, command: TradingCommand) -> anyhow::Result<()>
111    where
112        Self: ExecutionAlgorithmNative,
113        Self: 'static + std::fmt::Debug + Sized,
114    {
115        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
116        if core.config.log_commands {
117            let id = &core.actor.actor_id;
118            log::info!("{id} {RECV}{CMD} {command:?}");
119        }
120
121        if DataActorNative::core(core).state() != ComponentState::Running {
122            return Ok(());
123        }
124
125        match command {
126            TradingCommand::SubmitOrder(cmd) => {
127                self.subscribe_to_strategy_events(cmd.strategy_id);
128                let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
129                core.remember_submit_params(cmd.client_order_id, cmd.params.clone());
130                let order = core.get_order(&cmd.client_order_id)?;
131                self.on_order(order)
132            }
133            TradingCommand::SubmitOrderList(cmd) => {
134                self.subscribe_to_strategy_events(cmd.strategy_id);
135                let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
136                for client_order_id in &cmd.order_list.client_order_ids {
137                    core.remember_submit_params(*client_order_id, cmd.params.clone());
138                }
139                let orders = core.get_orders_for_list(&cmd.order_list)?;
140                self.on_order_list(cmd.order_list, orders)
141            }
142            TradingCommand::CancelOrder(cmd) => self.handle_cancel_order(cmd),
143            _ => {
144                log::warn!("Unhandled command type: {command:?}");
145                Ok(())
146            }
147        }
148    }
149
150    /// Called when a primary order is received for execution.
151    ///
152    /// Override this method to implement the algorithm's order slicing logic.
153    ///
154    /// # Errors
155    ///
156    /// Returns an error if order handling fails.
157    fn on_order(&mut self, order: OrderAny) -> anyhow::Result<()>;
158
159    /// Called when an order list is received for execution.
160    ///
161    /// Override this method to handle order lists. The default implementation
162    /// processes each order individually.
163    ///
164    /// # Errors
165    ///
166    /// Returns an error if order list handling fails.
167    fn on_order_list(
168        &mut self,
169        _order_list: OrderList,
170        orders: Vec<OrderAny>,
171    ) -> anyhow::Result<()> {
172        for order in orders {
173            self.on_order(order)?;
174        }
175        Ok(())
176    }
177
178    /// Denies an order by applying and publishing an `OrderDenied` event.
179    ///
180    /// An order absent from the cache is added first, with its `OrderInitialized` event published
181    /// before the denial. A closed cached order is left unchanged. Use an `OrderDeniedReason`
182    /// string for the standardized reason.
183    ///
184    /// # Errors
185    ///
186    /// Returns an error if:
187    /// - The algorithm is not registered with a trader.
188    /// - The order cannot be added to the cache.
189    /// - The denial cannot be applied, including an invalid order state transition.
190    ///
191    /// No event is published when the denial cannot be applied.
192    fn deny_order(&mut self, order: &OrderAny, reason: Ustr) -> anyhow::Result<()>
193    where
194        Self: ExecutionAlgorithmNative,
195    {
196        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
197        registered_trader_id(core)?;
198        let ts_now = core.clock_mut().timestamp_ns();
199        let event = OrderEventAny::Denied(OrderDenied::new(
200            order.trader_id(),
201            order.strategy_id(),
202            order.instrument_id(),
203            order.client_order_id(),
204            reason,
205            UUID4::new(),
206            ts_now,
207            ts_now,
208        ));
209
210        let publish_initialized = {
211            let cache_rc = core.cache_rc();
212            let mut cache = cache_rc.borrow_mut();
213
214            if cache
215                .order(&order.client_order_id())
216                .is_some_and(|cached_order| cached_order.is_closed())
217            {
218                return Ok(());
219            }
220
221            let publish_initialized = if cache.order_exists(&order.client_order_id()) {
222                false
223            } else {
224                cache.add_order(order.clone(), None, None, false)?;
225                true
226            };
227
228            cache.update_order(&event)?;
229            publish_initialized
230        };
231
232        if publish_initialized {
233            publish_order_initialized(order);
234        }
235        publish_order_event(&event);
236
237        // A denied order never executes, so its stored submit params are dropped here
238        // rather than waiting for an execution completion that will never arrive.
239        ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
240            .remove_submit_params(&order.client_order_id());
241
242        Ok(())
243    }
244
245    /// Handles a cancel order command for algorithm-managed orders.
246    ///
247    /// This generates an internal cancel event and publishes it. The order
248    /// is canceled locally without sending a command to the execution engine.
249    ///
250    /// # Errors
251    ///
252    /// Returns an error if cancellation fails.
253    fn handle_cancel_order(&mut self, command: CancelOrder) -> anyhow::Result<()>
254    where
255        Self: ExecutionAlgorithmNative,
256    {
257        let (order, is_pending_cancel) = {
258            let cache = ExecutionAlgorithmNative::exec_algorithm_core_mut(self).cache_ref();
259
260            let Some(order) = cache.order(&command.client_order_id) else {
261                log::warn!(
262                    "Cannot cancel order: {} not found in cache",
263                    command.client_order_id
264                );
265                return Ok(());
266            };
267
268            let is_pending = cache.is_order_pending_cancel_local(&command.client_order_id);
269            (order.clone(), is_pending)
270        };
271
272        if is_pending_cancel {
273            return Ok(());
274        }
275
276        if order.is_closed() {
277            log::warn!("Order already closed for {command:?}");
278            return Ok(());
279        }
280
281        let event = OrderEventAny::Canceled(self.generate_order_canceled(&order));
282
283        let order = {
284            let cache_rc = ExecutionAlgorithmNative::exec_algorithm_core_mut(self).cache_rc();
285            let mut cache = cache_rc.borrow_mut();
286            match cache.update_order(&event) {
287                Ok(order) => order,
288                Err(e)
289                    if matches!(
290                        e.downcast_ref::<OrderError>(),
291                        Some(OrderError::InvalidStateTransition)
292                    ) =>
293                {
294                    log::warn!("InvalidStateTrigger: {e}, did not apply cancel event");
295                    return Ok(());
296                }
297                Err(e) => return Err(e),
298            }
299        };
300
301        let topic = format!("events.order.{}", order.strategy_id());
302        msgbus::publish_order_event(topic.into(), &event);
303        msgbus::publish_order_event(
304            msgbus::switchboard::get_order_canceled_topic(order.instrument_id()),
305            &event,
306        );
307
308        Ok(())
309    }
310
311    /// Generates an `OrderCanceled` event for an order.
312    fn generate_order_canceled(&mut self, order: &OrderAny) -> OrderCanceled
313    where
314        Self: ExecutionAlgorithmNative,
315    {
316        let ts_now = ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
317            .clock_mut()
318            .timestamp_ns();
319
320        OrderCanceled::new(
321            order.trader_id(),
322            order.strategy_id(),
323            order.instrument_id(),
324            order.client_order_id(),
325            UUID4::new(),
326            ts_now,
327            ts_now,
328            false, // reconciliation
329            order.venue_order_id(),
330            order.account_id(),
331        )
332    }
333
334    /// Generates an `OrderPendingUpdate` event for an order.
335    fn generate_order_pending_update(&mut self, order: &OrderAny) -> OrderPendingUpdate
336    where
337        Self: ExecutionAlgorithmNative,
338    {
339        let ts_now = ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
340            .clock_mut()
341            .timestamp_ns();
342
343        OrderPendingUpdate::new(
344            order.trader_id(),
345            order.strategy_id(),
346            order.instrument_id(),
347            order.client_order_id(),
348            order.account_id(),
349            UUID4::new(),
350            ts_now,
351            ts_now,
352            false, // reconciliation
353            order.venue_order_id(),
354        )
355    }
356
357    /// Generates an `OrderPendingCancel` event for an order.
358    fn generate_order_pending_cancel(&mut self, order: &OrderAny) -> OrderPendingCancel
359    where
360        Self: ExecutionAlgorithmNative,
361    {
362        let ts_now = ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
363            .clock_mut()
364            .timestamp_ns();
365
366        OrderPendingCancel::new(
367            order.trader_id(),
368            order.strategy_id(),
369            order.instrument_id(),
370            order.client_order_id(),
371            order.account_id(),
372            UUID4::new(),
373            ts_now,
374            ts_now,
375            false, // reconciliation
376            order.venue_order_id(),
377        )
378    }
379
380    /// Spawns a market order from a primary order.
381    ///
382    /// Creates a new market order with:
383    /// - A unique client order ID: `{primary_id}-E{sequence}`.
384    /// - The primary order's trader ID, strategy ID, and instrument ID.
385    /// - The algorithm's `exec_algorithm_id`.
386    /// - `exec_spawn_id` set to the primary order's client order ID.
387    ///
388    /// If `reduce_primary` is true, the primary order's quantity will be reduced
389    /// by the spawned quantity. If the spawned order is subsequently denied or
390    /// rejected (before acceptance), the deducted quantity is automatically
391    /// restored to the primary order.
392    fn spawn_market(
393        &mut self,
394        primary: &mut OrderAny,
395        quantity: Quantity,
396        time_in_force: TimeInForce,
397        reduce_only: bool,
398        tags: Option<Vec<Ustr>>,
399        reduce_primary: bool,
400    ) -> MarketOrder
401    where
402        Self: ExecutionAlgorithmNative,
403    {
404        // Generate spawn ID first so we can track the reduction
405        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
406        let client_order_id = core.spawn_client_order_id(&primary.client_order_id());
407        let ts_init = core.clock_mut().timestamp_ns();
408        let exec_algorithm_id = core.exec_algorithm_id;
409
410        if reduce_primary {
411            self.reduce_primary_order(primary, quantity);
412            ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
413                .track_pending_spawn_reduction(client_order_id, quantity);
414        }
415
416        MarketOrder::new(
417            primary.trader_id(),
418            primary.strategy_id(),
419            primary.instrument_id(),
420            client_order_id,
421            primary.order_side(),
422            quantity,
423            time_in_force,
424            UUID4::new(),
425            ts_init,
426            reduce_only,
427            primary.is_quote_quantity(),
428            primary.contingency_type(),
429            primary.order_list_id(),
430            primary.linked_order_ids().map(|ids| ids.to_vec()),
431            primary.parent_order_id(),
432            Some(exec_algorithm_id),
433            primary.exec_algorithm_params().cloned(),
434            Some(primary.client_order_id()),
435            tags.or_else(|| primary.tags().map(|t| t.to_vec())),
436        )
437    }
438
439    /// Spawns a limit order from a primary order.
440    ///
441    /// Creates a new limit order with:
442    /// - A unique client order ID: `{primary_id}-E{sequence}`
443    /// - The primary order's trader ID, strategy ID, and instrument ID
444    /// - The algorithm's `exec_algorithm_id`
445    /// - `exec_spawn_id` set to the primary order's client order ID
446    ///
447    /// If `reduce_primary` is true, the primary order's quantity will be reduced
448    /// by the spawned quantity. If the spawned order is subsequently denied or
449    /// rejected (before acceptance), the deducted quantity is automatically
450    /// restored to the primary order.
451    #[expect(clippy::too_many_arguments)]
452    fn spawn_limit(
453        &mut self,
454        primary: &mut OrderAny,
455        quantity: Quantity,
456        price: Price,
457        time_in_force: TimeInForce,
458        expire_time: Option<UnixNanos>,
459        post_only: bool,
460        reduce_only: bool,
461        display_qty: Option<Quantity>,
462        emulation_trigger: Option<TriggerType>,
463        tags: Option<Vec<Ustr>>,
464        reduce_primary: bool,
465    ) -> LimitOrder
466    where
467        Self: ExecutionAlgorithmNative,
468    {
469        // Generate spawn ID first so we can track the reduction
470        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
471        let client_order_id = core.spawn_client_order_id(&primary.client_order_id());
472        let ts_init = core.clock_mut().timestamp_ns();
473        let exec_algorithm_id = core.exec_algorithm_id;
474
475        if reduce_primary {
476            self.reduce_primary_order(primary, quantity);
477            ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
478                .track_pending_spawn_reduction(client_order_id, quantity);
479        }
480
481        LimitOrder::new(
482            primary.trader_id(),
483            primary.strategy_id(),
484            primary.instrument_id(),
485            client_order_id,
486            primary.order_side(),
487            quantity,
488            price,
489            time_in_force,
490            expire_time,
491            post_only,
492            reduce_only,
493            primary.is_quote_quantity(),
494            display_qty,
495            emulation_trigger,
496            None, // trigger_instrument_id
497            primary.contingency_type(),
498            primary.order_list_id(),
499            primary.linked_order_ids().map(|ids| ids.to_vec()),
500            primary.parent_order_id(),
501            Some(exec_algorithm_id),
502            primary.exec_algorithm_params().cloned(),
503            Some(primary.client_order_id()),
504            tags.or_else(|| primary.tags().map(|t| t.to_vec())),
505            UUID4::new(),
506            ts_init,
507        )
508    }
509
510    /// Spawns a market-to-limit order from a primary order.
511    ///
512    /// Creates a new market-to-limit order with:
513    /// - A unique client order ID: `{primary_id}-E{sequence}`
514    /// - The primary order's trader ID, strategy ID, and instrument ID
515    /// - The algorithm's `exec_algorithm_id`
516    /// - `exec_spawn_id` set to the primary order's client order ID
517    ///
518    /// If `reduce_primary` is true, the primary order's quantity will be reduced
519    /// by the spawned quantity. If the spawned order is subsequently denied or
520    /// rejected (before acceptance), the deducted quantity is automatically
521    /// restored to the primary order.
522    #[expect(clippy::too_many_arguments)]
523    fn spawn_market_to_limit(
524        &mut self,
525        primary: &mut OrderAny,
526        quantity: Quantity,
527        time_in_force: TimeInForce,
528        expire_time: Option<UnixNanos>,
529        reduce_only: bool,
530        display_qty: Option<Quantity>,
531        emulation_trigger: Option<TriggerType>,
532        tags: Option<Vec<Ustr>>,
533        reduce_primary: bool,
534    ) -> MarketToLimitOrder
535    where
536        Self: ExecutionAlgorithmNative,
537    {
538        // Generate spawn ID first so we can track the reduction
539        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
540        let client_order_id = core.spawn_client_order_id(&primary.client_order_id());
541        let ts_init = core.clock_mut().timestamp_ns();
542        let exec_algorithm_id = core.exec_algorithm_id;
543
544        if reduce_primary {
545            self.reduce_primary_order(primary, quantity);
546            ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
547                .track_pending_spawn_reduction(client_order_id, quantity);
548        }
549
550        let mut order = MarketToLimitOrder::new(
551            primary.trader_id(),
552            primary.strategy_id(),
553            primary.instrument_id(),
554            client_order_id,
555            primary.order_side(),
556            quantity,
557            time_in_force,
558            expire_time,
559            false, // post_only
560            reduce_only,
561            primary.is_quote_quantity(),
562            display_qty,
563            primary.contingency_type(),
564            primary.order_list_id(),
565            primary.linked_order_ids().map(|ids| ids.to_vec()),
566            primary.parent_order_id(),
567            Some(exec_algorithm_id),
568            primary.exec_algorithm_params().cloned(),
569            Some(primary.client_order_id()),
570            tags.or_else(|| primary.tags().map(|t| t.to_vec())),
571            UUID4::new(),
572            ts_init,
573        );
574
575        if emulation_trigger.is_some() {
576            order.set_emulation_trigger(emulation_trigger);
577        }
578
579        order
580    }
581
582    /// Reduces the primary order's quantity by the spawn quantity.
583    ///
584    /// Generates an `OrderUpdated` event and applies it to the primary order,
585    /// then updates the order in the cache.
586    ///
587    /// # Panics
588    ///
589    /// Panics if `spawn_qty` exceeds the primary order's `leaves_qty`.
590    fn reduce_primary_order(&mut self, primary: &mut OrderAny, spawn_qty: Quantity)
591    where
592        Self: ExecutionAlgorithmNative,
593    {
594        let leaves_qty = primary.leaves_qty();
595        assert!(
596            leaves_qty >= spawn_qty,
597            "Spawn quantity {spawn_qty} exceeds primary leaves_qty {leaves_qty}"
598        );
599
600        let primary_qty = primary.quantity();
601        let new_qty = Quantity::from_raw(primary_qty.raw - spawn_qty.raw, primary_qty.precision);
602
603        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
604        let ts_now = core.clock_mut().timestamp_ns();
605
606        let updated = OrderUpdated::new(
607            primary.trader_id(),
608            primary.strategy_id(),
609            primary.instrument_id(),
610            primary.client_order_id(),
611            new_qty,
612            UUID4::new(),
613            ts_now,
614            ts_now,
615            false, // reconciliation
616            primary.venue_order_id(),
617            primary.account_id(),
618            None, // price
619            None, // trigger_price
620            None, // protection_price
621            primary.is_quote_quantity(),
622        );
623
624        let event = OrderEventAny::Updated(updated);
625
626        {
627            let cache_rc = core.cache_rc();
628            let mut cache = cache_rc.borrow_mut();
629            *primary = cache
630                .update_order(&event)
631                .expect("Failed to update order in cache");
632        }
633
634        publish_order_event(&event);
635    }
636
637    /// Restores the primary order quantity after a spawned order is denied or rejected.
638    ///
639    /// This is called when a spawned order fails before acceptance. The quantity
640    /// that was deducted from the primary order is restored (up to the spawned
641    /// order's `leaves_qty` to handle partial fills).
642    fn restore_primary_order_quantity(&mut self, order: &OrderAny)
643    where
644        Self: ExecutionAlgorithmNative,
645    {
646        let Some(exec_spawn_id) = order.exec_spawn_id() else {
647            return;
648        };
649
650        let reduction_qty = {
651            let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
652            core.take_pending_spawn_reduction(&order.client_order_id())
653        };
654
655        let Some(reduction_qty) = reduction_qty else {
656            return;
657        };
658
659        let primary = {
660            let cache = ExecutionAlgorithmNative::exec_algorithm_core_mut(self).cache_ref();
661            cache.order(&exec_spawn_id).map(|o| o.clone())
662        };
663
664        let Some(primary) = primary else {
665            log::warn!(
666                "Cannot restore primary order quantity: primary order {exec_spawn_id} not found",
667            );
668            return;
669        };
670
671        // Cap restore amount by leaves_qty to handle partial fills before rejection
672        let restore_raw = std::cmp::min(reduction_qty.raw, order.leaves_qty().raw);
673        if restore_raw == 0 {
674            return;
675        }
676
677        let restored_qty = Quantity::from_raw(
678            primary.quantity().raw + restore_raw,
679            primary.quantity().precision,
680        );
681
682        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
683        let ts_now = core.clock_mut().timestamp_ns();
684
685        let updated = OrderUpdated::new(
686            primary.trader_id(),
687            primary.strategy_id(),
688            primary.instrument_id(),
689            primary.client_order_id(),
690            restored_qty,
691            UUID4::new(),
692            ts_now,
693            ts_now,
694            false, // reconciliation
695            primary.venue_order_id(),
696            primary.account_id(),
697            None, // price
698            None, // trigger_price
699            None, // protection_price
700            primary.is_quote_quantity(),
701        );
702
703        let event = OrderEventAny::Updated(updated);
704
705        let primary = {
706            let cache_rc = core.cache_rc();
707            let mut cache = cache_rc.borrow_mut();
708            match cache.update_order(&event) {
709                Ok(primary) => primary,
710                Err(e) => {
711                    log::warn!("Failed to update primary order in cache: {e}");
712                    return;
713                }
714            }
715        };
716
717        publish_order_event(&event);
718
719        log::info!(
720            "Restored primary order {} quantity to {} after spawned order {} was denied/rejected",
721            primary.client_order_id(),
722            restored_qty,
723            order.client_order_id()
724        );
725    }
726
727    /// Submits an order to the execution engine via the risk engine.
728    ///
729    /// # Errors
730    ///
731    /// Returns an error if order submission fails.
732    fn submit_order(
733        &mut self,
734        order: OrderAny,
735        position_id: Option<PositionId>,
736        client_id: Option<ClientId>,
737    ) -> anyhow::Result<()>
738    where
739        Self: ExecutionAlgorithmNative,
740    {
741        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
742
743        let trader_id = registered_trader_id(core)?;
744        let ts_init = core.clock_mut().timestamp_ns();
745
746        // For spawned orders, use the parent's strategy ID
747        let strategy_id = order.strategy_id();
748
749        let primary_id = order
750            .exec_spawn_id()
751            .unwrap_or_else(|| order.client_order_id());
752        let params = core.submit_params(&primary_id);
753
754        let order_exists = {
755            let cache = core.cache_ref();
756            cache.order_exists(&order.client_order_id())
757        };
758
759        {
760            let cache_rc = core.cache_rc();
761            let mut cache = cache_rc.borrow_mut();
762            cache.add_order(order.clone(), position_id, client_id, true)?;
763        }
764
765        if !order_exists {
766            publish_order_initialized(&order);
767        }
768
769        let command = SubmitOrder::new(
770            trader_id,
771            client_id,
772            strategy_id,
773            order.instrument_id(),
774            order.client_order_id(),
775            order.init_event().clone(),
776            order.exec_algorithm_id(),
777            position_id,
778            params,
779            UUID4::new(),
780            ts_init,
781            None, // correlation_id
782        );
783
784        if core.config.log_commands {
785            let id = &core.actor.actor_id;
786            log::info!("{id} {SEND}{CMD} {command:?}");
787        }
788
789        msgbus::send_trading_command(
790            MessagingSwitchboard::risk_engine_execute(),
791            TradingCommand::SubmitOrder(command),
792        );
793
794        Ok(())
795    }
796
797    /// Modifies an order.
798    ///
799    /// # Errors
800    ///
801    /// Returns an error if order modification fails.
802    fn modify_order(
803        &mut self,
804        order: &mut OrderAny,
805        quantity: Option<Quantity>,
806        price: Option<Price>,
807        trigger_price: Option<Price>,
808        client_id: Option<ClientId>,
809    ) -> anyhow::Result<()>
810    where
811        Self: ExecutionAlgorithmNative,
812    {
813        let qty_changing = quantity.is_some_and(|q| q != order.quantity());
814        let price_changing = price.is_some() && price != order.price();
815        let trigger_changing = trigger_price.is_some() && trigger_price != order.trigger_price();
816
817        if !qty_changing && !price_changing && !trigger_changing {
818            log::error!(
819                "Cannot create command ModifyOrder: \
820                quantity, price, and trigger were either None \
821                or the same as existing values"
822            );
823            return Ok(());
824        }
825
826        if order.is_closed() || order.is_pending_cancel() {
827            log::warn!(
828                "Cannot create command ModifyOrder: state is {:?}, {order:?}",
829                order.status()
830            );
831            return Ok(());
832        }
833
834        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
835        let trader_id = registered_trader_id(core)?;
836        let strategy_id = order.strategy_id();
837
838        if !order.is_active_local() {
839            required_account_id(order, "pending update")?;
840            let event = self.generate_order_pending_update(order);
841            let event = OrderEventAny::PendingUpdate(event);
842
843            {
844                let cache_rc = ExecutionAlgorithmNative::exec_algorithm_core_mut(self).cache_rc();
845                let mut cache = cache_rc.borrow_mut();
846                match cache.update_order(&event) {
847                    Ok(updated) => *order = updated,
848                    Err(e)
849                        if matches!(
850                            e.downcast_ref::<OrderError>(),
851                            Some(OrderError::InvalidStateTransition)
852                        ) =>
853                    {
854                        log::warn!("InvalidStateTrigger: {e}, did not apply pending update event");
855                        return Ok(());
856                    }
857                    Err(e) => return Err(e),
858                }
859            }
860
861            let topic = format!("events.order.{strategy_id}");
862            msgbus::publish_order_event(topic.into(), &event);
863            msgbus::publish_order_event(
864                msgbus::switchboard::get_order_pending_update_topic(order.instrument_id()),
865                &event,
866            );
867        }
868
869        let ts_init = ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
870            .clock_mut()
871            .timestamp_ns();
872        let command = ModifyOrder::new(
873            trader_id,
874            client_id,
875            strategy_id,
876            order.instrument_id(),
877            order.client_order_id(),
878            order.venue_order_id(),
879            quantity,
880            price,
881            trigger_price,
882            UUID4::new(),
883            ts_init,
884            None, // params,
885            None, // correlation_id
886        );
887
888        if ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
889            .config
890            .log_commands
891        {
892            let id = &ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
893                .actor
894                .actor_id;
895            log::info!("{id} {SEND}{CMD} {command:?}");
896        }
897
898        let has_emulation_trigger = order
899            .emulation_trigger()
900            .is_some_and(|t| t != TriggerType::NoTrigger);
901
902        if order.is_emulated() || has_emulation_trigger {
903            msgbus::send_trading_command(
904                MessagingSwitchboard::order_emulator_execute(),
905                TradingCommand::ModifyOrder(command),
906            );
907        } else {
908            msgbus::send_trading_command(
909                MessagingSwitchboard::risk_engine_execute(),
910                TradingCommand::ModifyOrder(command),
911            );
912        }
913
914        Ok(())
915    }
916
917    /// Modifies an INITIALIZED or RELEASED order in place without sending a command.
918    ///
919    /// This is useful for adjusting order parameters before submission. The order
920    /// is updated locally by applying an `OrderUpdated` event and updating the cache.
921    ///
922    /// At least one parameter must differ from the current order values.
923    ///
924    /// # Errors
925    ///
926    /// Returns an error if the order status is not INITIALIZED or RELEASED,
927    /// or if no parameters would change.
928    fn modify_order_in_place(
929        &mut self,
930        order: &mut OrderAny,
931        quantity: Option<Quantity>,
932        price: Option<Price>,
933        trigger_price: Option<Price>,
934    ) -> anyhow::Result<()>
935    where
936        Self: ExecutionAlgorithmNative,
937    {
938        // Validate order status
939        let status = order.status();
940        if status != OrderStatus::Initialized && status != OrderStatus::Released {
941            anyhow::bail!(
942                "Cannot modify order in place: status is {status:?}, expected INITIALIZED or RELEASED"
943            );
944        }
945
946        // Validate order type compatibility
947        if price.is_some() && order.price().is_none() {
948            anyhow::bail!(
949                "Cannot modify order in place: {} orders do not have a LIMIT price",
950                order.order_type()
951            );
952        }
953
954        if trigger_price.is_some() && order.trigger_price().is_none() {
955            anyhow::bail!(
956                "Cannot modify order in place: {} orders do not have a STOP trigger price",
957                order.order_type()
958            );
959        }
960
961        // Check if any value would actually change
962        let qty_changing = quantity.is_some_and(|q| q != order.quantity());
963        let price_changing = price.is_some() && price != order.price();
964        let trigger_changing = trigger_price.is_some() && trigger_price != order.trigger_price();
965
966        if !qty_changing && !price_changing && !trigger_changing {
967            anyhow::bail!("Cannot modify order in place: no parameters differ from current values");
968        }
969
970        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
971        let ts_now = core.clock_mut().timestamp_ns();
972
973        let updated = OrderUpdated::new(
974            order.trader_id(),
975            order.strategy_id(),
976            order.instrument_id(),
977            order.client_order_id(),
978            quantity.unwrap_or_else(|| order.quantity()),
979            UUID4::new(),
980            ts_now,
981            ts_now,
982            false, // reconciliation
983            order.venue_order_id(),
984            order.account_id(),
985            price,
986            trigger_price,
987            None, // protection_price
988            order.is_quote_quantity(),
989        );
990
991        let event = OrderEventAny::Updated(updated);
992
993        {
994            let cache_rc = core.cache_rc();
995            let mut cache = cache_rc.borrow_mut();
996            *order = cache.update_order(&event)?;
997        }
998
999        publish_order_event(&event);
1000
1001        Ok(())
1002    }
1003
1004    /// Cancels an order.
1005    ///
1006    /// # Errors
1007    ///
1008    /// Returns an error if order cancellation fails.
1009    fn cancel_order(
1010        &mut self,
1011        order: &mut OrderAny,
1012        client_id: Option<ClientId>,
1013    ) -> anyhow::Result<()>
1014    where
1015        Self: ExecutionAlgorithmNative,
1016    {
1017        if order.is_closed() || order.is_pending_cancel() {
1018            log::warn!(
1019                "Cannot cancel order: state is {:?}, {order:?}",
1020                order.status()
1021            );
1022            return Ok(());
1023        }
1024
1025        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
1026        let trader_id = registered_trader_id(core)?;
1027        let strategy_id = order.strategy_id();
1028
1029        if !order.is_active_local() {
1030            required_account_id(order, "pending cancel")?;
1031            let event = self.generate_order_pending_cancel(order);
1032            let event = OrderEventAny::PendingCancel(event);
1033
1034            {
1035                let cache_rc = ExecutionAlgorithmNative::exec_algorithm_core_mut(self).cache_rc();
1036                let mut cache = cache_rc.borrow_mut();
1037                match cache.update_order(&event) {
1038                    Ok(updated) => *order = updated,
1039                    Err(e)
1040                        if matches!(
1041                            e.downcast_ref::<OrderError>(),
1042                            Some(OrderError::InvalidStateTransition)
1043                        ) =>
1044                    {
1045                        log::warn!("InvalidStateTrigger: {e}, did not apply pending cancel event");
1046                        return Ok(());
1047                    }
1048                    Err(e) => return Err(e),
1049                }
1050            }
1051
1052            let topic = format!("events.order.{strategy_id}");
1053            msgbus::publish_order_event(topic.into(), &event);
1054            msgbus::publish_order_event(
1055                msgbus::switchboard::get_order_pending_cancel_topic(order.instrument_id()),
1056                &event,
1057            );
1058        }
1059
1060        let ts_init = ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
1061            .clock_mut()
1062            .timestamp_ns();
1063        let command = CancelOrder::new(
1064            trader_id,
1065            client_id,
1066            strategy_id,
1067            order.instrument_id(),
1068            order.client_order_id(),
1069            order.venue_order_id(),
1070            UUID4::new(),
1071            ts_init,
1072            None, // params,
1073            None, // correlation_id
1074        );
1075
1076        if ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
1077            .config
1078            .log_commands
1079        {
1080            let id = &ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
1081                .actor
1082                .actor_id;
1083            log::info!("{id} {SEND}{CMD} {command:?}");
1084        }
1085
1086        let has_emulation_trigger = order
1087            .emulation_trigger()
1088            .is_some_and(|t| t != TriggerType::NoTrigger);
1089
1090        if order.is_emulated() || order.status() == OrderStatus::Released || has_emulation_trigger {
1091            msgbus::send_trading_command(
1092                MessagingSwitchboard::order_emulator_execute(),
1093                TradingCommand::CancelOrder(command),
1094            );
1095        } else {
1096            msgbus::send_trading_command(
1097                MessagingSwitchboard::exec_engine_execute(),
1098                TradingCommand::CancelOrder(command),
1099            );
1100        }
1101
1102        Ok(())
1103    }
1104
1105    /// Subscribes to events from a strategy.
1106    ///
1107    /// This is called automatically when the first order is received from a strategy.
1108    fn subscribe_to_strategy_events(&mut self, strategy_id: StrategyId)
1109    where
1110        Self: ExecutionAlgorithmNative,
1111        Self: 'static + std::fmt::Debug + Sized,
1112    {
1113        let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
1114        if core.is_strategy_subscribed(&strategy_id) {
1115            return;
1116        }
1117
1118        let actor_id = core.actor.actor_id.inner();
1119
1120        let order_topic = format!("events.order.{strategy_id}");
1121        let order_actor_id = actor_id;
1122        let order_handler = TypedHandler::from(move |event: &OrderEventAny| {
1123            if let Some(mut algo) = try_get_actor_unchecked::<Self>(&order_actor_id) {
1124                algo.handle_order_event(event.clone());
1125            } else {
1126                log::error!(
1127                    "ExecutionAlgorithm {order_actor_id} not found for order event handling"
1128                );
1129            }
1130        });
1131        msgbus::subscribe_order_events(order_topic.clone().into(), order_handler.clone(), None);
1132
1133        let position_topic = format!("events.position.{strategy_id}");
1134        let position_handler = TypedHandler::from(move |event: &PositionEvent| {
1135            if let Some(mut algo) = try_get_actor_unchecked::<Self>(&actor_id) {
1136                algo.handle_position_event(event.clone());
1137            } else {
1138                log::error!("ExecutionAlgorithm {actor_id} not found for position event handling");
1139            }
1140        });
1141        msgbus::subscribe_position_events(
1142            position_topic.clone().into(),
1143            position_handler.clone(),
1144            None,
1145        );
1146
1147        let handlers = StrategyEventHandlers {
1148            order_topic,
1149            order_handler,
1150            position_topic,
1151            position_handler,
1152        };
1153        core.store_strategy_event_handlers(strategy_id, handlers);
1154
1155        core.add_subscribed_strategy(strategy_id);
1156        log::info!("Subscribed to events for strategy {strategy_id}");
1157    }
1158
1159    /// Unsubscribes from all strategy event handlers.
1160    ///
1161    /// This should be called before reset to properly clean up msgbus subscriptions.
1162    fn unsubscribe_all_strategy_events(&mut self)
1163    where
1164        Self: ExecutionAlgorithmNative,
1165    {
1166        let handlers =
1167            ExecutionAlgorithmNative::exec_algorithm_core_mut(self).take_strategy_event_handlers();
1168
1169        for (strategy_id, h) in handlers {
1170            msgbus::unsubscribe_order_events(h.order_topic.into(), &h.order_handler);
1171            msgbus::unsubscribe_position_events(h.position_topic.into(), &h.position_handler);
1172            log::info!("Unsubscribed from events for strategy {strategy_id}");
1173        }
1174        ExecutionAlgorithmNative::exec_algorithm_core_mut(self).clear_subscribed_strategies();
1175    }
1176
1177    /// Handles an order event, filtering for algorithm-owned orders.
1178    fn handle_order_event(&mut self, event: OrderEventAny)
1179    where
1180        Self: ExecutionAlgorithmNative,
1181    {
1182        if DataActorNative::core(ExecutionAlgorithmNative::exec_algorithm_core_mut(self)).state()
1183            != ComponentState::Running
1184        {
1185            return;
1186        }
1187
1188        let order = {
1189            let cache = ExecutionAlgorithmNative::exec_algorithm_core_mut(self).cache_ref();
1190            cache.order(&event.client_order_id()).map(|o| o.clone())
1191        };
1192
1193        let Some(order) = order else {
1194            return;
1195        };
1196
1197        let Some(order_algo_id) = order.exec_algorithm_id() else {
1198            return;
1199        };
1200
1201        if order_algo_id != self.id() {
1202            return;
1203        }
1204
1205        {
1206            let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
1207            if core.config.log_events {
1208                let id = &core.actor.actor_id;
1209                log::info!("{id} {RECV}{EVT} {event}");
1210            }
1211        }
1212
1213        match &event {
1214            OrderEventAny::Initialized(e) => self.on_order_initialized(e.clone()),
1215            OrderEventAny::Denied(e) => {
1216                self.restore_primary_order_quantity(&order);
1217                self.on_order_denied(*e);
1218            }
1219            OrderEventAny::Emulated(e) => self.on_order_emulated(*e),
1220            OrderEventAny::Released(e) => self.on_order_released(*e),
1221            OrderEventAny::Submitted(e) => self.on_order_submitted(*e),
1222            OrderEventAny::Rejected(e) => {
1223                self.restore_primary_order_quantity(&order);
1224                self.on_order_rejected(*e);
1225            }
1226            OrderEventAny::Accepted(e) => {
1227                // Commit reduction - order accepted by venue
1228                ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
1229                    .take_pending_spawn_reduction(&order.client_order_id());
1230                self.on_order_accepted(*e);
1231            }
1232            OrderEventAny::Canceled(e) => {
1233                ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
1234                    .take_pending_spawn_reduction(&order.client_order_id());
1235                self.on_algo_order_canceled(*e);
1236            }
1237            OrderEventAny::Expired(e) => {
1238                ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
1239                    .take_pending_spawn_reduction(&order.client_order_id());
1240                self.on_order_expired(*e);
1241            }
1242            OrderEventAny::Triggered(e) => self.on_order_triggered(*e),
1243            OrderEventAny::PendingUpdate(e) => self.on_order_pending_update(*e),
1244            OrderEventAny::PendingCancel(e) => self.on_order_pending_cancel(*e),
1245            OrderEventAny::ModifyRejected(e) => self.on_order_modify_rejected(*e),
1246            OrderEventAny::CancelRejected(e) => self.on_order_cancel_rejected(*e),
1247            OrderEventAny::Updated(e) => self.on_order_updated(*e),
1248            OrderEventAny::Filled(e) => self.on_algo_order_filled(e.clone()),
1249            OrderEventAny::FillVoided(e) => self.on_order_fill_voided(e),
1250        }
1251
1252        self.on_order_event(event);
1253    }
1254
1255    /// Handles a position event.
1256    fn handle_position_event(&mut self, event: PositionEvent)
1257    where
1258        Self: ExecutionAlgorithmNative,
1259    {
1260        if DataActorNative::core(ExecutionAlgorithmNative::exec_algorithm_core_mut(self)).state()
1261            != ComponentState::Running
1262        {
1263            return;
1264        }
1265
1266        {
1267            let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
1268            if core.config.log_events {
1269                let id = &core.actor.actor_id;
1270                log::info!("{id} {RECV}{EVT} {event:?}");
1271            }
1272        }
1273
1274        match &event {
1275            PositionEvent::PositionOpened(e) => self.on_position_opened(e.clone()),
1276            PositionEvent::PositionChanged(e) => self.on_position_changed(e.clone()),
1277            PositionEvent::PositionClosed(e) => self.on_position_closed(e.clone()),
1278            PositionEvent::PositionAdjusted(_) => {}
1279        }
1280
1281        self.on_position_event(event);
1282    }
1283
1284    /// Called when the algorithm is started.
1285    ///
1286    /// Override this method to implement custom initialization logic.
1287    ///
1288    /// # Errors
1289    ///
1290    /// Returns an error if start fails.
1291    fn on_start(&mut self) -> anyhow::Result<()>
1292    where
1293        Self: ExecutionAlgorithmNative,
1294    {
1295        let id = self.id();
1296        log::info!("Starting {id}");
1297        Ok(())
1298    }
1299
1300    /// Called when the algorithm is stopped.
1301    ///
1302    /// # Errors
1303    ///
1304    /// Returns an error if stop fails.
1305    fn on_stop(&mut self) -> anyhow::Result<()> {
1306        Ok(())
1307    }
1308
1309    /// Called when the algorithm is reset.
1310    ///
1311    /// # Errors
1312    ///
1313    /// Returns an error if reset fails.
1314    fn on_reset(&mut self) -> anyhow::Result<()>
1315    where
1316        Self: ExecutionAlgorithmNative,
1317    {
1318        self.unsubscribe_all_strategy_events();
1319        ExecutionAlgorithmNative::exec_algorithm_core_mut(self).reset();
1320        Ok(())
1321    }
1322
1323    /// Called when a time event is received.
1324    ///
1325    /// Override this method for timer-based algorithms like TWAP.
1326    ///
1327    /// # Errors
1328    ///
1329    /// Returns an error if time event handling fails.
1330    fn on_time_event(&mut self, _event: &TimeEvent) -> anyhow::Result<()> {
1331        Ok(())
1332    }
1333
1334    /// Called when an order is initialized.
1335    #[allow(unused_variables)]
1336    fn on_order_initialized(&mut self, event: OrderInitialized) {}
1337
1338    /// Called when an order is denied.
1339    #[allow(unused_variables)]
1340    fn on_order_denied(&mut self, event: OrderDenied) {}
1341
1342    /// Called when an order is emulated.
1343    #[allow(unused_variables)]
1344    fn on_order_emulated(&mut self, event: OrderEmulated) {}
1345
1346    /// Called when an order is released from emulation.
1347    #[allow(unused_variables)]
1348    fn on_order_released(&mut self, event: OrderReleased) {}
1349
1350    /// Called when an order is submitted.
1351    #[allow(unused_variables)]
1352    fn on_order_submitted(&mut self, event: OrderSubmitted) {}
1353
1354    /// Called when an order is rejected.
1355    #[allow(unused_variables)]
1356    fn on_order_rejected(&mut self, event: OrderRejected) {}
1357
1358    /// Called when an order is accepted.
1359    #[allow(unused_variables)]
1360    fn on_order_accepted(&mut self, event: OrderAccepted) {}
1361
1362    /// Called when an order is canceled.
1363    #[allow(unused_variables)]
1364    fn on_algo_order_canceled(&mut self, event: OrderCanceled) {}
1365
1366    /// Called when an order expires.
1367    #[allow(unused_variables)]
1368    fn on_order_expired(&mut self, event: OrderExpired) {}
1369
1370    /// Called when an order is triggered.
1371    #[allow(unused_variables)]
1372    fn on_order_triggered(&mut self, event: OrderTriggered) {}
1373
1374    /// Called when an order modification is pending.
1375    #[allow(unused_variables)]
1376    fn on_order_pending_update(&mut self, event: OrderPendingUpdate) {}
1377
1378    /// Called when an order cancellation is pending.
1379    #[allow(unused_variables)]
1380    fn on_order_pending_cancel(&mut self, event: OrderPendingCancel) {}
1381
1382    /// Called when an order modification is rejected.
1383    #[allow(unused_variables)]
1384    fn on_order_modify_rejected(&mut self, event: OrderModifyRejected) {}
1385
1386    /// Called when an order cancellation is rejected.
1387    #[allow(unused_variables)]
1388    fn on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {}
1389
1390    /// Called when an order is updated.
1391    #[allow(unused_variables)]
1392    fn on_order_updated(&mut self, event: OrderUpdated) {}
1393
1394    /// Called when an order is filled.
1395    #[allow(unused_variables)]
1396    fn on_algo_order_filled(&mut self, event: OrderFilled) {}
1397
1398    /// Called when an applied order fill is partly or fully voided.
1399    #[allow(unused_variables)]
1400    fn on_order_fill_voided(&mut self, event: &OrderFillVoided) {}
1401
1402    /// Called for any order event (after specific handler).
1403    #[allow(unused_variables)]
1404    fn on_order_event(&mut self, event: OrderEventAny) {}
1405
1406    /// Called when a position is opened.
1407    #[allow(unused_variables)]
1408    fn on_position_opened(&mut self, event: PositionOpened) {}
1409
1410    /// Called when a position is changed.
1411    #[allow(unused_variables)]
1412    fn on_position_changed(&mut self, event: PositionChanged) {}
1413
1414    /// Called when a position is closed.
1415    #[allow(unused_variables)]
1416    fn on_position_closed(&mut self, event: PositionClosed) {}
1417
1418    /// Called for any position event (after specific handler).
1419    #[allow(unused_variables)]
1420    fn on_position_event(&mut self, event: PositionEvent) {}
1421}
1422
1423fn publish_order_initialized(order: &OrderAny) {
1424    let event = OrderEventAny::Initialized(order.init_event().clone());
1425    publish_order_event(&event);
1426}
1427
1428fn publish_order_event(event: &OrderEventAny) {
1429    let topic = format!("events.order.{}", event.strategy_id());
1430    msgbus::publish_order_event(topic.into(), event);
1431}
1432
1433fn registered_trader_id(core: &ExecutionAlgorithmCore) -> anyhow::Result<TraderId> {
1434    DataActorNative::core(core)
1435        .trader_id()
1436        .ok_or_else(|| anyhow::anyhow!("ExecutionAlgorithm not registered: trader_id is not set"))
1437}
1438
1439fn required_account_id(order: &OrderAny, operation: &str) -> anyhow::Result<AccountId> {
1440    order.account_id().ok_or_else(|| {
1441        anyhow::anyhow!(
1442            "Cannot generate {operation} event for {}: account_id is not set",
1443            order.client_order_id()
1444        )
1445    })
1446}
1447
1448#[cfg(test)]
1449mod tests {
1450    use std::{cell::RefCell, rc::Rc};
1451
1452    use nautilus_common::{
1453        actor::DataActor,
1454        cache::Cache,
1455        clock::{Clock, TestClock},
1456        component::Component,
1457        enums::ComponentTrigger,
1458        msgbus,
1459        msgbus::TypedHandler,
1460    };
1461    use nautilus_model::{
1462        enums::{OrderSide, OrderStatus, OrderType},
1463        events::{
1464            OrderAccepted, OrderCanceled, OrderDenied, OrderDeniedReason, OrderRejected,
1465            order::spec::{
1466                OrderAcceptedSpec, OrderCanceledSpec, OrderDeniedSpec, OrderFillVoidedSpec,
1467                OrderFilledSpec, OrderRejectedSpec,
1468            },
1469        },
1470        identifiers::{
1471            AccountId, ActorId, ClientOrderId, ComponentId, ExecAlgorithmId, InstrumentId,
1472            StrategyId, TraderId, VenueOrderId,
1473        },
1474        orders::{LimitOrder, MarketOrder, OrderAny, OrderTestBuilder, stubs::TestOrderStubs},
1475        types::{Price, Quantity},
1476    };
1477    use rstest::rstest;
1478
1479    use super::*;
1480    use crate::nautilus_execution_algorithm;
1481
1482    #[derive(Debug)]
1483    struct TestAlgorithm {
1484        core: ExecutionAlgorithmCore,
1485        order_client_ids: Vec<ClientOrderId>,
1486    }
1487
1488    #[derive(Debug)]
1489    struct CoreFreeExecutionAlgorithm {
1490        state: ComponentState,
1491        orders_seen: usize,
1492    }
1493
1494    #[derive(Debug)]
1495    struct MacroTestCustomField {
1496        inner: ExecutionAlgorithmCore,
1497    }
1498
1499    impl Component for CoreFreeExecutionAlgorithm {
1500        fn component_id(&self) -> ComponentId {
1501            ComponentId::new("CoreFreeExecutionAlgorithm")
1502        }
1503
1504        fn state(&self) -> ComponentState {
1505            self.state
1506        }
1507
1508        fn transition_state(&mut self, trigger: ComponentTrigger) -> anyhow::Result<()> {
1509            self.state = self.state.transition(&trigger)?;
1510            Ok(())
1511        }
1512
1513        fn register(
1514            &mut self,
1515            _trader_id: TraderId,
1516            _clock: Rc<RefCell<dyn Clock>>,
1517            _cache: Rc<RefCell<Cache>>,
1518        ) -> anyhow::Result<()> {
1519            Ok(())
1520        }
1521    }
1522
1523    impl DataActor for CoreFreeExecutionAlgorithm {}
1524
1525    impl ExecutionAlgorithm for CoreFreeExecutionAlgorithm {
1526        fn on_order(&mut self, _order: OrderAny) -> anyhow::Result<()> {
1527            self.orders_seen += 1;
1528            Ok(())
1529        }
1530    }
1531
1532    impl DataActor for MacroTestCustomField {}
1533
1534    nautilus_execution_algorithm!(MacroTestCustomField, inner, {
1535        fn on_order(&mut self, _order: OrderAny) -> anyhow::Result<()> {
1536            Ok(())
1537        }
1538    });
1539
1540    impl TestAlgorithm {
1541        fn new(config: ExecutionAlgorithmConfig) -> Self {
1542            Self {
1543                core: ExecutionAlgorithmCore::new(config),
1544                order_client_ids: Vec::new(),
1545            }
1546        }
1547    }
1548
1549    impl DataActor for TestAlgorithm {}
1550
1551    nautilus_execution_algorithm!(TestAlgorithm, {
1552        fn on_order(&mut self, order: OrderAny) -> anyhow::Result<()> {
1553            self.order_client_ids.push(order.client_order_id());
1554            Ok(())
1555        }
1556    });
1557
1558    fn create_test_algorithm() -> TestAlgorithm {
1559        // Use unique ID to avoid thread-local registry/msgbus conflicts in parallel tests
1560        let unique_id = format!("TEST-{}", UUID4::new());
1561        let config = ExecutionAlgorithmConfig {
1562            exec_algorithm_id: Some(ExecAlgorithmId::new(&unique_id)),
1563            ..Default::default()
1564        };
1565        TestAlgorithm::new(config)
1566    }
1567
1568    fn register_algorithm(algo: &mut TestAlgorithm) {
1569        let trader_id = TraderId::from("TRADER-001");
1570        let clock = Rc::new(RefCell::new(TestClock::new()));
1571        let cache = Rc::new(RefCell::new(Cache::default()));
1572
1573        algo.core.register(trader_id, clock, cache).unwrap();
1574
1575        // Transition to Running state for tests
1576        algo.transition_state(ComponentTrigger::Initialize).unwrap();
1577        algo.transition_state(ComponentTrigger::Start).unwrap();
1578        algo.transition_state(ComponentTrigger::StartCompleted)
1579            .unwrap();
1580    }
1581
1582    fn subscribe_order_topic(
1583        strategy_id: StrategyId,
1584    ) -> (TypedHandler<OrderEventAny>, Rc<RefCell<Vec<OrderEventAny>>>) {
1585        let events = Rc::new(RefCell::new(Vec::new()));
1586        let handler = TypedHandler::from({
1587            let events = events.clone();
1588            move |event: &OrderEventAny| {
1589                events.borrow_mut().push(event.clone());
1590            }
1591        });
1592        msgbus::subscribe_order_events(
1593            format!("events.order.{strategy_id}").into(),
1594            handler.clone(),
1595            None,
1596        );
1597        (handler, events)
1598    }
1599
1600    #[rstest]
1601    fn test_algorithm_creation() {
1602        let algo = create_test_algorithm();
1603        assert!(algo.id().inner().starts_with("TEST-"));
1604        assert!(algo.order_client_ids.is_empty());
1605    }
1606
1607    #[rstest]
1608    fn test_algorithm_registration() {
1609        let mut algo = create_test_algorithm();
1610        register_algorithm(&mut algo);
1611
1612        assert_eq!(algo.trader_id(), Some(TraderId::from("TRADER-001")));
1613    }
1614
1615    #[rstest]
1616    fn test_algorithm_deny_order_updates_cache_and_publishes_once() {
1617        let mut algo = create_test_algorithm();
1618        register_algorithm(&mut algo);
1619
1620        let strategy_id = StrategyId::from("STRAT-ALGO-DENY");
1621        let order = OrderAny::Market(MarketOrder::new(
1622            TraderId::from("TRADER-001"),
1623            strategy_id,
1624            InstrumentId::from("BTC/USDT.BINANCE"),
1625            ClientOrderId::from("O-ALGO-DENY"),
1626            OrderSide::Buy,
1627            Quantity::from("1.0"),
1628            TimeInForce::Gtc,
1629            UUID4::new(),
1630            0.into(),
1631            false,
1632            false,
1633            None,
1634            None,
1635            None,
1636            None,
1637            None,
1638            None,
1639            None,
1640            None,
1641        ));
1642        {
1643            let cache_rc = algo.core.cache_rc();
1644            cache_rc
1645                .borrow_mut()
1646                .add_order(order.clone(), None, None, false)
1647                .unwrap();
1648        }
1649        let reason = OrderDeniedReason::ValidationFailed {
1650            detail: "invalid execution schedule".to_string(),
1651        }
1652        .to_string();
1653        let reason = Ustr::from(&reason);
1654        let (handler, events) = subscribe_order_topic(strategy_id);
1655
1656        algo.deny_order(&order, reason).unwrap();
1657        algo.deny_order(&order, reason).unwrap();
1658
1659        msgbus::unsubscribe_order_events(format!("events.order.{strategy_id}").into(), &handler);
1660        let cached_order = algo.cache().order(&order.client_order_id()).unwrap();
1661        let events = events.borrow();
1662
1663        assert_eq!(cached_order.status(), OrderStatus::Denied);
1664        assert_eq!(events.len(), 1);
1665        assert!(matches!(
1666            &events[0],
1667            OrderEventAny::Denied(event)
1668                if event.reason == reason
1669                    && event.strategy_id == strategy_id
1670                    && event.client_order_id == order.client_order_id()
1671        ));
1672    }
1673
1674    #[rstest]
1675    fn test_algorithm_deny_order_initializes_missing_order_once() {
1676        let mut algo = create_test_algorithm();
1677        register_algorithm(&mut algo);
1678
1679        let strategy_id = StrategyId::from("STRAT-ALGO-DENY-MISSING");
1680        let order = OrderTestBuilder::new(OrderType::Market)
1681            .strategy_id(strategy_id)
1682            .instrument_id(InstrumentId::from("BTC/USDT.BINANCE"))
1683            .client_order_id(ClientOrderId::from("O-ALGO-DENY-MISSING"))
1684            .quantity(Quantity::from("1.0"))
1685            .build();
1686        let reason = Ustr::from("VALIDATION_FAILED: invalid execution schedule");
1687        let (handler, events) = subscribe_order_topic(strategy_id);
1688
1689        algo.deny_order(&order, reason).unwrap();
1690        algo.deny_order(&order, reason).unwrap();
1691
1692        msgbus::unsubscribe_order_events(format!("events.order.{strategy_id}").into(), &handler);
1693        let cached_order = algo.cache().order(&order.client_order_id()).unwrap();
1694        let events = events.borrow();
1695
1696        assert_eq!(cached_order.status(), OrderStatus::Denied);
1697        assert_eq!(cached_order.event_count(), 2);
1698        assert_eq!(events.len(), 2);
1699        assert!(matches!(
1700            &events[0],
1701            OrderEventAny::Initialized(event)
1702                if event.strategy_id == strategy_id
1703                    && event.client_order_id == order.client_order_id()
1704        ));
1705        assert!(matches!(
1706            &events[1],
1707            OrderEventAny::Denied(event)
1708                if event.reason == reason
1709                    && event.strategy_id == strategy_id
1710                    && event.client_order_id == order.client_order_id()
1711        ));
1712    }
1713
1714    #[rstest]
1715    fn test_algorithm_deny_order_does_not_publish_when_apply_fails() {
1716        let mut algo = create_test_algorithm();
1717        register_algorithm(&mut algo);
1718
1719        let strategy_id = StrategyId::from("STRAT-ALGO-DENY-APPLY");
1720        let order = OrderTestBuilder::new(OrderType::Market)
1721            .strategy_id(strategy_id)
1722            .instrument_id(InstrumentId::from("BTC/USDT.BINANCE"))
1723            .client_order_id(ClientOrderId::from("O-ALGO-DENY-APPLY"))
1724            .quantity(Quantity::from("1.0"))
1725            .build();
1726        let order = TestOrderStubs::make_accepted_order(&order);
1727        {
1728            let cache_rc = algo.core.cache_rc();
1729            cache_rc
1730                .borrow_mut()
1731                .add_order(order.clone(), None, None, false)
1732                .unwrap();
1733        }
1734        let (handler, events) = subscribe_order_topic(strategy_id);
1735
1736        let mut params = nautilus_core::Params::new();
1737        params.insert(
1738            "route".to_string(),
1739            serde_json::Value::String("A".to_string()),
1740        );
1741        algo.core
1742            .remember_submit_params(order.client_order_id(), Some(params));
1743
1744        let error = algo
1745            .deny_order(
1746                &order,
1747                Ustr::from("VALIDATION_FAILED: invalid execution schedule"),
1748            )
1749            .unwrap_err();
1750
1751        msgbus::unsubscribe_order_events(format!("events.order.{strategy_id}").into(), &handler);
1752        let cached_order = algo.cache().order(&order.client_order_id()).unwrap();
1753
1754        assert!(matches!(
1755            error.downcast_ref::<OrderError>(),
1756            Some(OrderError::InvalidStateTransition)
1757        ));
1758        assert_eq!(cached_order.status(), OrderStatus::Accepted);
1759        assert!(events.borrow().is_empty());
1760        // A failed denial is not terminal, so the submit params must be retained
1761        assert!(algo.core.submit_params(&order.client_order_id()).is_some());
1762    }
1763
1764    #[rstest]
1765    fn test_algorithm_deny_order_removes_submit_params() {
1766        let mut algo = create_test_algorithm();
1767        register_algorithm(&mut algo);
1768
1769        let strategy_id = StrategyId::from("STRAT-ALGO-DENY-PARAMS");
1770        let order = OrderTestBuilder::new(OrderType::Market)
1771            .strategy_id(strategy_id)
1772            .instrument_id(InstrumentId::from("BTC/USDT.BINANCE"))
1773            .client_order_id(ClientOrderId::from("O-ALGO-DENY-PARAMS"))
1774            .quantity(Quantity::from("1.0"))
1775            .build();
1776        {
1777            let cache_rc = algo.core.cache_rc();
1778            cache_rc
1779                .borrow_mut()
1780                .add_order(order.clone(), None, None, false)
1781                .unwrap();
1782        }
1783
1784        let mut params = nautilus_core::Params::new();
1785        params.insert(
1786            "route".to_string(),
1787            serde_json::Value::String("A".to_string()),
1788        );
1789        algo.core
1790            .remember_submit_params(order.client_order_id(), Some(params));
1791        assert!(algo.core.submit_params(&order.client_order_id()).is_some());
1792
1793        algo.deny_order(&order, Ustr::from("VALIDATION_FAILED: test"))
1794            .unwrap();
1795
1796        assert!(algo.core.submit_params(&order.client_order_id()).is_none());
1797    }
1798
1799    #[rstest]
1800    fn test_submit_order_errors_when_algorithm_not_registered() {
1801        let mut algo = create_test_algorithm();
1802        let order = OrderAny::Market(MarketOrder::new(
1803            TraderId::from("TRADER-001"),
1804            StrategyId::from("STRAT-001"),
1805            InstrumentId::from("BTC/USDT.BINANCE"),
1806            ClientOrderId::from("O-UNREGISTERED-001"),
1807            OrderSide::Buy,
1808            Quantity::from("1.0"),
1809            TimeInForce::Gtc,
1810            UUID4::new(),
1811            0.into(),
1812            false,
1813            false,
1814            None,
1815            None,
1816            None,
1817            None,
1818            None,
1819            None,
1820            None,
1821            None,
1822        ));
1823
1824        let err = algo
1825            .submit_order(order, None, None)
1826            .unwrap_err()
1827            .to_string();
1828
1829        assert_eq!(
1830            err,
1831            "ExecutionAlgorithm not registered: trader_id is not set"
1832        );
1833    }
1834
1835    #[rstest]
1836    fn test_required_account_id_errors_when_missing_for_algorithm_event() {
1837        let order = OrderAny::Market(MarketOrder::new(
1838            TraderId::from("TRADER-001"),
1839            StrategyId::from("STRAT-001"),
1840            InstrumentId::from("BTC/USDT.BINANCE"),
1841            ClientOrderId::from("O-NO-ACCOUNT-001"),
1842            OrderSide::Buy,
1843            Quantity::from("1.0"),
1844            TimeInForce::Gtc,
1845            UUID4::new(),
1846            0.into(),
1847            false,
1848            false,
1849            None,
1850            None,
1851            None,
1852            None,
1853            None,
1854            None,
1855            None,
1856            None,
1857        ));
1858
1859        let err = required_account_id(&order, "pending update")
1860            .unwrap_err()
1861            .to_string();
1862
1863        assert_eq!(
1864            err,
1865            "Cannot generate pending update event for O-NO-ACCOUNT-001: account_id is not set"
1866        );
1867    }
1868
1869    #[rstest]
1870    fn test_algorithm_id() {
1871        let algo = create_test_algorithm();
1872        assert!(algo.id().inner().starts_with("TEST-"));
1873    }
1874
1875    #[rstest]
1876    fn test_execution_algorithm_behavior_does_not_require_native_core_access() {
1877        fn assert_execution_algorithm<T: ExecutionAlgorithm + DataActor + Component>() {}
1878
1879        assert_execution_algorithm::<CoreFreeExecutionAlgorithm>();
1880
1881        let mut algorithm = CoreFreeExecutionAlgorithm {
1882            state: ComponentState::PreInitialized,
1883            orders_seen: 0,
1884        };
1885        let order = OrderTestBuilder::new(OrderType::Market)
1886            .instrument_id(InstrumentId::from("BTC/USDT.BINANCE"))
1887            .quantity(Quantity::from("1.0"))
1888            .build();
1889
1890        algorithm.on_order(order).unwrap();
1891
1892        assert_eq!(algorithm.orders_seen, 1);
1893    }
1894
1895    #[rstest]
1896    fn test_nautilus_execution_algorithm_macro_custom_field() {
1897        let exec_algorithm_id = ExecAlgorithmId::from("MACRO-001");
1898        let algorithm = MacroTestCustomField {
1899            inner: ExecutionAlgorithmCore::new(ExecutionAlgorithmConfig {
1900                exec_algorithm_id: Some(exec_algorithm_id),
1901                ..Default::default()
1902            }),
1903        };
1904
1905        assert_eq!(algorithm.id(), exec_algorithm_id);
1906        assert_eq!(algorithm.actor_id(), ActorId::from("MACRO-001"));
1907    }
1908
1909    #[rstest]
1910    fn test_algorithm_spawn_market_creates_valid_order() {
1911        let mut algo = create_test_algorithm();
1912        register_algorithm(&mut algo);
1913
1914        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
1915        let mut primary = OrderAny::Market(MarketOrder::new(
1916            TraderId::from("TRADER-001"),
1917            StrategyId::from("STRAT-001"),
1918            instrument_id,
1919            ClientOrderId::from("O-001"),
1920            OrderSide::Buy,
1921            Quantity::from("1.0"),
1922            TimeInForce::Gtc,
1923            UUID4::new(),
1924            0.into(),
1925            false, // reduce_only
1926            false, // quote_quantity
1927            None,  // contingency_type
1928            None,  // order_list_id
1929            None,  // linked_order_ids
1930            None,  // parent_order_id
1931            None,  // exec_algorithm_id
1932            None,  // exec_algorithm_params
1933            None,  // exec_spawn_id
1934            None,  // tags
1935        ));
1936
1937        let spawned = algo.spawn_market(
1938            &mut primary,
1939            Quantity::from("0.5"),
1940            TimeInForce::Ioc,
1941            false,
1942            None,  // tags
1943            false, // reduce_primary
1944        );
1945
1946        assert_eq!(spawned.client_order_id.as_str(), "O-001-E1");
1947        assert_eq!(spawned.instrument_id, instrument_id);
1948        assert_eq!(spawned.order_side(), OrderSide::Buy);
1949        assert_eq!(spawned.quantity, Quantity::from("0.5"));
1950        assert_eq!(spawned.time_in_force, TimeInForce::Ioc);
1951        assert_eq!(spawned.exec_algorithm_id, Some(algo.id()));
1952        assert_eq!(spawned.exec_spawn_id, Some(ClientOrderId::from("O-001")));
1953    }
1954
1955    #[rstest]
1956    fn test_algorithm_spawn_increments_sequence() {
1957        let mut algo = create_test_algorithm();
1958        register_algorithm(&mut algo);
1959
1960        let mut primary = OrderAny::Market(MarketOrder::new(
1961            TraderId::from("TRADER-001"),
1962            StrategyId::from("STRAT-001"),
1963            InstrumentId::from("BTC/USDT.BINANCE"),
1964            ClientOrderId::from("O-001"),
1965            OrderSide::Buy,
1966            Quantity::from("1.0"),
1967            TimeInForce::Gtc,
1968            UUID4::new(),
1969            0.into(),
1970            false,
1971            false,
1972            None,
1973            None,
1974            None,
1975            None,
1976            None,
1977            None,
1978            None,
1979            None,
1980        ));
1981
1982        let spawned1 = algo.spawn_market(
1983            &mut primary,
1984            Quantity::from("0.25"),
1985            TimeInForce::Ioc,
1986            false,
1987            None,
1988            false,
1989        );
1990        let spawned2 = algo.spawn_market(
1991            &mut primary,
1992            Quantity::from("0.25"),
1993            TimeInForce::Ioc,
1994            false,
1995            None,
1996            false,
1997        );
1998        let spawned3 = algo.spawn_market(
1999            &mut primary,
2000            Quantity::from("0.25"),
2001            TimeInForce::Ioc,
2002            false,
2003            None,
2004            false,
2005        );
2006
2007        assert_eq!(spawned1.client_order_id.as_str(), "O-001-E1");
2008        assert_eq!(spawned2.client_order_id.as_str(), "O-001-E2");
2009        assert_eq!(spawned3.client_order_id.as_str(), "O-001-E3");
2010    }
2011
2012    #[rstest]
2013    fn test_algorithm_default_handlers_do_not_panic() {
2014        let mut algo = create_test_algorithm();
2015
2016        algo.on_order_initialized(OrderInitialized::default());
2017        algo.on_order_denied(OrderDenied::default());
2018        algo.on_order_emulated(OrderEmulated::default());
2019        algo.on_order_released(OrderReleased::default());
2020        algo.on_order_submitted(OrderSubmitted::default());
2021        algo.on_order_rejected(OrderRejected::default());
2022        algo.on_order_accepted(OrderAccepted::default());
2023        algo.on_algo_order_canceled(OrderCanceled::default());
2024        algo.on_order_expired(OrderExpired::default());
2025        algo.on_order_triggered(OrderTriggered::default());
2026        algo.on_order_pending_update(OrderPendingUpdate::default());
2027        algo.on_order_pending_cancel(OrderPendingCancel::default());
2028        algo.on_order_modify_rejected(OrderModifyRejected::default());
2029        algo.on_order_cancel_rejected(OrderCancelRejected::default());
2030        algo.on_order_updated(OrderUpdated::default());
2031        algo.on_algo_order_filled(OrderFilledSpec::builder().build());
2032        algo.on_order_fill_voided(&OrderFillVoidedSpec::builder().build());
2033    }
2034
2035    #[rstest]
2036    fn test_strategy_subscription_tracking() {
2037        let mut algo = create_test_algorithm();
2038        let strategy_id = StrategyId::from("STRAT-001");
2039
2040        assert!(!algo.core.is_strategy_subscribed(&strategy_id));
2041
2042        algo.subscribe_to_strategy_events(strategy_id);
2043        assert!(algo.core.is_strategy_subscribed(&strategy_id));
2044
2045        // Second call should be idempotent
2046        algo.subscribe_to_strategy_events(strategy_id);
2047        assert!(algo.core.is_strategy_subscribed(&strategy_id));
2048    }
2049
2050    #[rstest]
2051    fn test_algorithm_reset() {
2052        let mut algo = create_test_algorithm();
2053        let strategy_id = StrategyId::from("STRAT-001");
2054        let primary_id = ClientOrderId::new("O-001");
2055
2056        let _ = algo.core.spawn_client_order_id(&primary_id);
2057        algo.core.add_subscribed_strategy(strategy_id);
2058
2059        assert!(algo.core.spawn_sequence(&primary_id).is_some());
2060        assert!(algo.core.is_strategy_subscribed(&strategy_id));
2061
2062        ExecutionAlgorithm::on_reset(&mut algo).unwrap();
2063
2064        assert!(algo.core.spawn_sequence(&primary_id).is_none());
2065        assert!(!algo.core.is_strategy_subscribed(&strategy_id));
2066    }
2067
2068    #[rstest]
2069    fn test_algorithm_spawn_limit_creates_valid_order() {
2070        let mut algo = create_test_algorithm();
2071        register_algorithm(&mut algo);
2072
2073        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
2074        let mut primary = OrderAny::Market(MarketOrder::new(
2075            TraderId::from("TRADER-001"),
2076            StrategyId::from("STRAT-001"),
2077            instrument_id,
2078            ClientOrderId::from("O-001"),
2079            OrderSide::Buy,
2080            Quantity::from("1.0"),
2081            TimeInForce::Gtc,
2082            UUID4::new(),
2083            0.into(),
2084            false,
2085            false,
2086            None,
2087            None,
2088            None,
2089            None,
2090            None,
2091            None,
2092            None,
2093            None,
2094        ));
2095
2096        let price = Price::from("50000.0");
2097        let spawned = algo.spawn_limit(
2098            &mut primary,
2099            Quantity::from("0.5"),
2100            price,
2101            TimeInForce::Gtc,
2102            None,  // expire_time
2103            false, // post_only
2104            false, // reduce_only
2105            None,  // display_qty
2106            None,  // emulation_trigger
2107            None,  // tags
2108            false, // reduce_primary
2109        );
2110
2111        assert_eq!(spawned.client_order_id.as_str(), "O-001-E1");
2112        assert_eq!(spawned.instrument_id, instrument_id);
2113        assert_eq!(spawned.order_side(), OrderSide::Buy);
2114        assert_eq!(spawned.quantity, Quantity::from("0.5"));
2115        assert_eq!(spawned.price, price);
2116        assert_eq!(spawned.time_in_force, TimeInForce::Gtc);
2117        assert_eq!(spawned.exec_algorithm_id, Some(algo.id()));
2118        assert_eq!(spawned.exec_spawn_id, Some(ClientOrderId::from("O-001")));
2119    }
2120
2121    #[rstest]
2122    fn test_algorithm_spawn_market_to_limit_creates_valid_order() {
2123        let mut algo = create_test_algorithm();
2124        register_algorithm(&mut algo);
2125
2126        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
2127        let mut primary = OrderAny::Market(MarketOrder::new(
2128            TraderId::from("TRADER-001"),
2129            StrategyId::from("STRAT-001"),
2130            instrument_id,
2131            ClientOrderId::from("O-001"),
2132            OrderSide::Buy,
2133            Quantity::from("1.0"),
2134            TimeInForce::Gtc,
2135            UUID4::new(),
2136            0.into(),
2137            false,
2138            false,
2139            None,
2140            None,
2141            None,
2142            None,
2143            None,
2144            None,
2145            None,
2146            None,
2147        ));
2148
2149        let spawned = algo.spawn_market_to_limit(
2150            &mut primary,
2151            Quantity::from("0.5"),
2152            TimeInForce::Gtc,
2153            None,  // expire_time
2154            false, // reduce_only
2155            None,  // display_qty
2156            None,  // emulation_trigger
2157            None,  // tags
2158            false, // reduce_primary
2159        );
2160
2161        assert_eq!(spawned.client_order_id.as_str(), "O-001-E1");
2162        assert_eq!(spawned.instrument_id, instrument_id);
2163        assert_eq!(spawned.order_side(), OrderSide::Buy);
2164        assert_eq!(spawned.quantity, Quantity::from("0.5"));
2165        assert_eq!(spawned.time_in_force, TimeInForce::Gtc);
2166        assert_eq!(spawned.exec_algorithm_id, Some(algo.id()));
2167        assert_eq!(spawned.exec_spawn_id, Some(ClientOrderId::from("O-001")));
2168    }
2169
2170    #[rstest]
2171    fn test_algorithm_spawn_market_with_tags() {
2172        let mut algo = create_test_algorithm();
2173        register_algorithm(&mut algo);
2174
2175        let mut primary = OrderAny::Market(MarketOrder::new(
2176            TraderId::from("TRADER-001"),
2177            StrategyId::from("STRAT-001"),
2178            InstrumentId::from("BTC/USDT.BINANCE"),
2179            ClientOrderId::from("O-001"),
2180            OrderSide::Buy,
2181            Quantity::from("1.0"),
2182            TimeInForce::Gtc,
2183            UUID4::new(),
2184            0.into(),
2185            false,
2186            false,
2187            None,
2188            None,
2189            None,
2190            None,
2191            None,
2192            None,
2193            None,
2194            None,
2195        ));
2196
2197        let tags = vec![ustr::Ustr::from("TAG1"), ustr::Ustr::from("TAG2")];
2198        let spawned = algo.spawn_market(
2199            &mut primary,
2200            Quantity::from("0.5"),
2201            TimeInForce::Ioc,
2202            false,
2203            Some(tags.clone()),
2204            false,
2205        );
2206
2207        assert_eq!(spawned.tags, Some(tags));
2208    }
2209
2210    #[rstest]
2211    fn test_algorithm_spawn_propagates_primary_fields() {
2212        let mut algo = create_test_algorithm();
2213        register_algorithm(&mut algo);
2214
2215        let mut params = indexmap::IndexMap::new();
2216        params.insert(ustr::Ustr::from("horizon_secs"), ustr::Ustr::from("30"));
2217        params.insert(ustr::Ustr::from("interval_secs"), ustr::Ustr::from("10"));
2218        let primary_tags = vec![ustr::Ustr::from("PRIMARY_TAG")];
2219        let linked_order_ids = vec![ClientOrderId::from("LINK-1")];
2220        let client_order_id = ClientOrderId::from("O-001");
2221
2222        let mut primary = OrderAny::Market(MarketOrder::new(
2223            TraderId::from("TRADER-001"),
2224            StrategyId::from("STRAT-001"),
2225            InstrumentId::from("BTC/USDT.BINANCE"),
2226            client_order_id,
2227            OrderSide::Buy,
2228            Quantity::from("1.0"),
2229            TimeInForce::Gtc,
2230            UUID4::new(),
2231            0.into(),
2232            false, // reduce_only
2233            true,  // quote_quantity
2234            None,  // contingency_type
2235            None,  // order_list_id
2236            Some(linked_order_ids.clone()),
2237            None, // parent_order_id
2238            Some(algo.id()),
2239            Some(params.clone()),
2240            Some(client_order_id),
2241            Some(primary_tags.clone()),
2242        ));
2243
2244        let spawned_market = algo.spawn_market(
2245            &mut primary,
2246            Quantity::from("0.25"),
2247            TimeInForce::Ioc,
2248            false,
2249            None, // falls back to primary.tags
2250            false,
2251        );
2252        assert!(spawned_market.is_quote_quantity);
2253        assert_eq!(spawned_market.exec_algorithm_params, Some(params.clone()));
2254        assert_eq!(spawned_market.tags, Some(primary_tags.clone()));
2255        assert_eq!(
2256            spawned_market.linked_order_ids,
2257            Some(linked_order_ids.clone())
2258        );
2259
2260        let spawned_limit = algo.spawn_limit(
2261            &mut primary,
2262            Quantity::from("0.25"),
2263            Price::from("50000.0"),
2264            TimeInForce::Gtc,
2265            None,  // expire_time
2266            false, // post_only
2267            false, // reduce_only
2268            None,  // display_qty
2269            None,  // emulation_trigger
2270            None,  // falls back to primary.tags
2271            false,
2272        );
2273        assert!(spawned_limit.is_quote_quantity);
2274        assert_eq!(spawned_limit.exec_algorithm_params, Some(params.clone()));
2275        assert_eq!(spawned_limit.tags, Some(primary_tags.clone()));
2276        assert_eq!(
2277            spawned_limit.linked_order_ids,
2278            Some(linked_order_ids.clone())
2279        );
2280
2281        let spawned_mtl = algo.spawn_market_to_limit(
2282            &mut primary,
2283            Quantity::from("0.25"),
2284            TimeInForce::Gtc,
2285            None,  // expire_time
2286            false, // reduce_only
2287            None,  // display_qty
2288            None,  // emulation_trigger
2289            None,  // falls back to primary.tags
2290            false,
2291        );
2292        assert!(spawned_mtl.is_quote_quantity);
2293        assert_eq!(spawned_mtl.exec_algorithm_params, Some(params));
2294        assert_eq!(spawned_mtl.tags, Some(primary_tags));
2295        assert_eq!(spawned_mtl.linked_order_ids, Some(linked_order_ids));
2296    }
2297
2298    #[rstest]
2299    fn test_algorithm_reduce_primary_order() {
2300        let mut algo = create_test_algorithm();
2301        register_algorithm(&mut algo);
2302
2303        let order = OrderAny::Market(MarketOrder::new(
2304            TraderId::from("TRADER-001"),
2305            StrategyId::from("STRAT-001"),
2306            InstrumentId::from("BTC/USDT.BINANCE"),
2307            ClientOrderId::from("O-001"),
2308            OrderSide::Buy,
2309            Quantity::from("1.0"),
2310            TimeInForce::Gtc,
2311            UUID4::new(),
2312            0.into(),
2313            false,
2314            false,
2315            None,
2316            None,
2317            None,
2318            None,
2319            None,
2320            None,
2321            None,
2322            None,
2323        ));
2324
2325        // Make accepted so OrderUpdated can be applied
2326        let mut primary = TestOrderStubs::make_accepted_order(&order);
2327
2328        {
2329            let cache_rc = algo.core.cache_rc();
2330            let mut cache = cache_rc.borrow_mut();
2331            cache.add_order(primary.clone(), None, None, false).unwrap();
2332        }
2333
2334        let spawn_qty = Quantity::from("0.3");
2335        algo.reduce_primary_order(&mut primary, spawn_qty);
2336
2337        assert_eq!(primary.quantity(), Quantity::from("0.7"));
2338    }
2339
2340    #[rstest]
2341    fn test_algorithm_reduce_primary_order_publishes_updated_event() {
2342        let mut algo = create_test_algorithm();
2343        register_algorithm(&mut algo);
2344
2345        let strategy_id = StrategyId::from("STRAT-ALGO-REDUCE-PUBLISH");
2346        let order = OrderAny::Market(MarketOrder::new(
2347            TraderId::from("TRADER-001"),
2348            strategy_id,
2349            InstrumentId::from("BTC/USDT.BINANCE"),
2350            ClientOrderId::from("O-ALGO-REDUCE"),
2351            OrderSide::Buy,
2352            Quantity::from("1.0"),
2353            TimeInForce::Gtc,
2354            UUID4::new(),
2355            0.into(),
2356            false,
2357            false,
2358            None,
2359            None,
2360            None,
2361            None,
2362            None,
2363            None,
2364            None,
2365            None,
2366        ));
2367        let mut primary = TestOrderStubs::make_accepted_order(&order);
2368
2369        {
2370            let cache_rc = algo.core.cache_rc();
2371            let mut cache = cache_rc.borrow_mut();
2372            cache.add_order(primary.clone(), None, None, false).unwrap();
2373        }
2374
2375        let (handler, events) = subscribe_order_topic(strategy_id);
2376
2377        algo.reduce_primary_order(&mut primary, Quantity::from("0.3"));
2378
2379        msgbus::unsubscribe_order_events(format!("events.order.{strategy_id}").into(), &handler);
2380        let events = events.borrow();
2381
2382        assert_eq!(events.len(), 1);
2383        assert!(matches!(
2384            &events[0],
2385            OrderEventAny::Updated(event) if event.quantity == Quantity::from("0.7")
2386        ));
2387    }
2388
2389    #[rstest]
2390    fn test_algorithm_submit_order_publishes_initialized_for_new_order() {
2391        let mut algo = create_test_algorithm();
2392        register_algorithm(&mut algo);
2393
2394        let strategy_id = StrategyId::from("STRAT-ALGO-INIT-PUBLISH");
2395        let order = OrderAny::Market(MarketOrder::new(
2396            TraderId::from("TRADER-001"),
2397            strategy_id,
2398            InstrumentId::from("BTC/USDT.BINANCE"),
2399            ClientOrderId::from("O-ALGO-INIT"),
2400            OrderSide::Buy,
2401            Quantity::from("1.0"),
2402            TimeInForce::Gtc,
2403            UUID4::new(),
2404            0.into(),
2405            false,
2406            false,
2407            None,
2408            None,
2409            None,
2410            None,
2411            None,
2412            None,
2413            None,
2414            None,
2415        ));
2416        let (handler, events) = subscribe_order_topic(strategy_id);
2417
2418        algo.submit_order(order.clone(), None, None).unwrap();
2419
2420        msgbus::unsubscribe_order_events(format!("events.order.{strategy_id}").into(), &handler);
2421        let events = events.borrow();
2422
2423        assert_eq!(events.len(), 1);
2424        assert!(matches!(
2425            &events[0],
2426            OrderEventAny::Initialized(event) if event.client_order_id == order.client_order_id()
2427        ));
2428    }
2429
2430    #[rstest]
2431    fn test_algorithm_submit_order_does_not_republish_initialized_for_existing_order() {
2432        let mut algo = create_test_algorithm();
2433        register_algorithm(&mut algo);
2434
2435        let strategy_id = StrategyId::from("STRAT-ALGO-INIT-EXISTING");
2436        let order = OrderAny::Market(MarketOrder::new(
2437            TraderId::from("TRADER-001"),
2438            strategy_id,
2439            InstrumentId::from("BTC/USDT.BINANCE"),
2440            ClientOrderId::from("O-ALGO-INIT-EXISTING"),
2441            OrderSide::Buy,
2442            Quantity::from("1.0"),
2443            TimeInForce::Gtc,
2444            UUID4::new(),
2445            0.into(),
2446            false,
2447            false,
2448            None,
2449            None,
2450            None,
2451            None,
2452            None,
2453            None,
2454            None,
2455            None,
2456        ));
2457        {
2458            let cache_rc = algo.core.cache_rc();
2459            let mut cache = cache_rc.borrow_mut();
2460            cache.add_order(order.clone(), None, None, true).unwrap();
2461        }
2462        let (handler, events) = subscribe_order_topic(strategy_id);
2463
2464        algo.submit_order(order, None, None).unwrap();
2465
2466        msgbus::unsubscribe_order_events(format!("events.order.{strategy_id}").into(), &handler);
2467        assert!(events.borrow().is_empty());
2468    }
2469
2470    #[rstest]
2471    fn test_algorithm_spawn_market_with_reduce_primary() {
2472        let mut algo = create_test_algorithm();
2473        register_algorithm(&mut algo);
2474
2475        let order = OrderAny::Market(MarketOrder::new(
2476            TraderId::from("TRADER-001"),
2477            StrategyId::from("STRAT-001"),
2478            InstrumentId::from("BTC/USDT.BINANCE"),
2479            ClientOrderId::from("O-001"),
2480            OrderSide::Buy,
2481            Quantity::from("1.0"),
2482            TimeInForce::Gtc,
2483            UUID4::new(),
2484            0.into(),
2485            false,
2486            false,
2487            None,
2488            None,
2489            None,
2490            None,
2491            None,
2492            None,
2493            None,
2494            None,
2495        ));
2496
2497        // Make accepted so OrderUpdated can be applied
2498        let mut primary = TestOrderStubs::make_accepted_order(&order);
2499
2500        {
2501            let cache_rc = algo.core.cache_rc();
2502            let mut cache = cache_rc.borrow_mut();
2503            cache.add_order(primary.clone(), None, None, false).unwrap();
2504        }
2505
2506        let spawned = algo.spawn_market(
2507            &mut primary,
2508            Quantity::from("0.4"),
2509            TimeInForce::Ioc,
2510            false,
2511            None,
2512            true, // reduce_primary = true
2513        );
2514
2515        assert_eq!(spawned.quantity, Quantity::from("0.4"));
2516        assert_eq!(primary.quantity(), Quantity::from("0.6"));
2517    }
2518    #[rstest]
2519    fn test_algorithm_forwards_captured_params_to_spawned_order() {
2520        let mut algo = create_test_algorithm();
2521        register_algorithm(&mut algo);
2522
2523        let strategy_id = StrategyId::from("STRAT-FWD-001");
2524        let mut primary = OrderAny::Market(MarketOrder::new(
2525            TraderId::from("TRADER-001"),
2526            strategy_id,
2527            InstrumentId::from("BTC/USDT.BINANCE"),
2528            ClientOrderId::from("O-FWD-001"),
2529            OrderSide::Buy,
2530            Quantity::from("1.0"),
2531            TimeInForce::Gtc,
2532            UUID4::new(),
2533            0.into(),
2534            false,
2535            false,
2536            None,
2537            None,
2538            None,
2539            None,
2540            None,
2541            None,
2542            None,
2543            None,
2544        ));
2545        {
2546            let cache_rc = algo.core.cache_rc();
2547            let mut cache = cache_rc.borrow_mut();
2548            cache.add_order(primary.clone(), None, None, true).unwrap();
2549        }
2550
2551        let mut params = nautilus_core::Params::new();
2552        params.insert("is_leverage".to_string(), serde_json::Value::Bool(true));
2553        let command = SubmitOrder::new(
2554            TraderId::from("TRADER-001"),
2555            None,
2556            strategy_id,
2557            primary.instrument_id(),
2558            primary.client_order_id(),
2559            primary.init_event().clone(),
2560            primary.exec_algorithm_id(),
2561            None,
2562            Some(params),
2563            UUID4::new(),
2564            0.into(),
2565            None,
2566        );
2567        algo.execute(TradingCommand::SubmitOrder(command)).unwrap();
2568
2569        let received = Rc::new(RefCell::new(None::<SubmitOrder>));
2570        let handler = msgbus::TypedIntoHandler::from({
2571            let captured = received.clone();
2572            move |cmd: TradingCommand| {
2573                if let TradingCommand::SubmitOrder(cmd) = cmd {
2574                    *captured.borrow_mut() = Some(cmd);
2575                }
2576            }
2577        });
2578        msgbus::register_trading_command_endpoint(
2579            MessagingSwitchboard::risk_engine_execute(),
2580            handler,
2581        );
2582
2583        let spawned = algo.spawn_market(
2584            &mut primary,
2585            Quantity::from("0.4"),
2586            TimeInForce::Ioc,
2587            false,
2588            None,
2589            false, // reduce_primary
2590        );
2591        algo.submit_order(OrderAny::Market(spawned), None, None)
2592            .unwrap();
2593
2594        let captured = received.borrow();
2595        let cmd = captured.as_ref().expect("expected a forwarded SubmitOrder");
2596        assert_eq!(cmd.client_order_id, ClientOrderId::from("O-FWD-001-E1"));
2597        assert_eq!(
2598            cmd.params.as_ref().and_then(|p| p.get_bool("is_leverage")),
2599            Some(true),
2600        );
2601    }
2602
2603    #[rstest]
2604    fn test_algorithm_submit_order_list_captures_params_per_order() {
2605        use nautilus_common::messages::execution::SubmitOrderList;
2606        use nautilus_model::identifiers::OrderListId;
2607
2608        let mut algo = create_test_algorithm();
2609        register_algorithm(&mut algo);
2610
2611        let strategy_id = StrategyId::from("STRAT-LIST-001");
2612        let order1 = OrderAny::Market(MarketOrder::new(
2613            TraderId::from("TRADER-001"),
2614            strategy_id,
2615            InstrumentId::from("BTC/USDT.BINANCE"),
2616            ClientOrderId::from("O-LIST-001"),
2617            OrderSide::Buy,
2618            Quantity::from("1.0"),
2619            TimeInForce::Gtc,
2620            UUID4::new(),
2621            0.into(),
2622            false,
2623            false,
2624            None,
2625            None,
2626            None,
2627            None,
2628            None,
2629            None,
2630            None,
2631            None,
2632        ));
2633        let order2 = OrderAny::Market(MarketOrder::new(
2634            TraderId::from("TRADER-001"),
2635            strategy_id,
2636            InstrumentId::from("BTC/USDT.BINANCE"),
2637            ClientOrderId::from("O-LIST-002"),
2638            OrderSide::Buy,
2639            Quantity::from("1.0"),
2640            TimeInForce::Gtc,
2641            UUID4::new(),
2642            0.into(),
2643            false,
2644            false,
2645            None,
2646            None,
2647            None,
2648            None,
2649            None,
2650            None,
2651            None,
2652            None,
2653        ));
2654        {
2655            let cache_rc = algo.core.cache_rc();
2656            let mut cache = cache_rc.borrow_mut();
2657            cache.add_order(order1.clone(), None, None, true).unwrap();
2658            cache.add_order(order2.clone(), None, None, true).unwrap();
2659        }
2660
2661        let order_list = OrderList::new(
2662            OrderListId::from("OL-001"),
2663            order1.instrument_id(),
2664            strategy_id,
2665            vec![order1.client_order_id(), order2.client_order_id()],
2666            0.into(),
2667        );
2668
2669        let mut params = nautilus_core::Params::new();
2670        params.insert("is_leverage".to_string(), serde_json::Value::Bool(true));
2671        let command = SubmitOrderList::new(
2672            TraderId::from("TRADER-001"),
2673            None,
2674            strategy_id,
2675            order_list,
2676            vec![order1.init_event().clone(), order2.init_event().clone()],
2677            order1.exec_algorithm_id(),
2678            None,
2679            Some(params),
2680            UUID4::new(),
2681            0.into(),
2682            None,
2683        );
2684        algo.execute(TradingCommand::SubmitOrderList(command))
2685            .unwrap();
2686
2687        assert_eq!(
2688            algo.order_client_ids,
2689            [
2690                ClientOrderId::from("O-LIST-001"),
2691                ClientOrderId::from("O-LIST-002"),
2692            ],
2693        );
2694
2695        for id in ["O-LIST-001", "O-LIST-002"] {
2696            assert_eq!(
2697                algo.core
2698                    .submit_params(&ClientOrderId::from(id))
2699                    .and_then(|p| p.get_bool("is_leverage")),
2700                Some(true),
2701                "expected forwarded params for {id}",
2702            );
2703        }
2704    }
2705
2706    #[rstest]
2707    fn test_algorithm_generate_order_canceled() {
2708        let mut algo = create_test_algorithm();
2709        register_algorithm(&mut algo);
2710
2711        let order = OrderAny::Market(MarketOrder::new(
2712            TraderId::from("TRADER-001"),
2713            StrategyId::from("STRAT-001"),
2714            InstrumentId::from("BTC/USDT.BINANCE"),
2715            ClientOrderId::from("O-001"),
2716            OrderSide::Buy,
2717            Quantity::from("1.0"),
2718            TimeInForce::Gtc,
2719            UUID4::new(),
2720            0.into(),
2721            false,
2722            false,
2723            None,
2724            None,
2725            None,
2726            None,
2727            None,
2728            None,
2729            None,
2730            None,
2731        ));
2732
2733        let event = algo.generate_order_canceled(&order);
2734
2735        assert_eq!(event.trader_id, TraderId::from("TRADER-001"));
2736        assert_eq!(event.strategy_id, StrategyId::from("STRAT-001"));
2737        assert_eq!(event.instrument_id, InstrumentId::from("BTC/USDT.BINANCE"));
2738        assert_eq!(event.client_order_id, ClientOrderId::from("O-001"));
2739    }
2740
2741    #[rstest]
2742    fn test_algorithm_handle_cancel_order_publishes_instrument_canceled_topic() {
2743        let mut algo = create_test_algorithm();
2744        register_algorithm(&mut algo);
2745
2746        let strategy_id = StrategyId::from("STRAT-ALGO-CANCEL-PUBLISH");
2747        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
2748        let order = OrderAny::Market(MarketOrder::new(
2749            TraderId::from("TRADER-001"),
2750            strategy_id,
2751            instrument_id,
2752            ClientOrderId::from("O-ALGO-CANCEL"),
2753            OrderSide::Buy,
2754            Quantity::from("1.0"),
2755            TimeInForce::Gtc,
2756            UUID4::new(),
2757            0.into(),
2758            false,
2759            false,
2760            None,
2761            None,
2762            None,
2763            None,
2764            None,
2765            None,
2766            None,
2767            None,
2768        ));
2769        let order = TestOrderStubs::make_accepted_order(&order);
2770
2771        {
2772            let cache_rc = algo.core.cache_rc();
2773            let mut cache = cache_rc.borrow_mut();
2774            cache.add_order(order.clone(), None, None, false).unwrap();
2775        }
2776
2777        let received = Rc::new(RefCell::new(Vec::<OrderEventAny>::new()));
2778        let handler = TypedHandler::from({
2779            let received = received.clone();
2780            move |event: &OrderEventAny| {
2781                received.borrow_mut().push(event.clone());
2782            }
2783        });
2784        let topic = msgbus::switchboard::get_order_canceled_topic(instrument_id);
2785        msgbus::subscribe_order_events(topic.into(), handler.clone(), None);
2786
2787        let command = CancelOrder::new(
2788            order.trader_id(),
2789            None,
2790            strategy_id,
2791            instrument_id,
2792            order.client_order_id(),
2793            order.venue_order_id(),
2794            UUID4::new(),
2795            0.into(),
2796            None,
2797            None,
2798        );
2799        algo.handle_cancel_order(command).unwrap();
2800
2801        msgbus::unsubscribe_order_events(topic.into(), &handler);
2802        let received = received.borrow();
2803        assert_eq!(received.len(), 1);
2804        assert!(matches!(received[0], OrderEventAny::Canceled(_)));
2805        assert_eq!(received[0].client_order_id(), order.client_order_id());
2806        assert_eq!(received[0].instrument_id(), instrument_id);
2807    }
2808
2809    #[rstest]
2810    fn test_algorithm_modify_order_in_place_updates_quantity() {
2811        let mut algo = create_test_algorithm();
2812        register_algorithm(&mut algo);
2813
2814        let strategy_id = StrategyId::from("STRAT-ALGO-MODIFY-IN-PLACE");
2815        let mut order = OrderAny::Limit(LimitOrder::new(
2816            TraderId::from("TRADER-001"),
2817            strategy_id,
2818            InstrumentId::from("BTC/USDT.BINANCE"),
2819            ClientOrderId::from("O-001"),
2820            OrderSide::Buy,
2821            Quantity::from("1.0"),
2822            Price::from("50000.0"),
2823            TimeInForce::Gtc,
2824            None,  // expire_time
2825            false, // post_only
2826            false, // reduce_only
2827            false, // quote_quantity
2828            None,  // display_qty
2829            None,  // emulation_trigger
2830            None,  // trigger_instrument_id
2831            None,  // contingency_type
2832            None,  // order_list_id
2833            None,  // linked_order_ids
2834            None,  // parent_order_id
2835            None,  // exec_algorithm_id
2836            None,  // exec_algorithm_params
2837            None,  // exec_spawn_id
2838            None,  // tags
2839            UUID4::new(),
2840            0.into(),
2841        ));
2842
2843        {
2844            let cache_rc = algo.core.cache_rc();
2845            let mut cache = cache_rc.borrow_mut();
2846            cache.add_order(order.clone(), None, None, false).unwrap();
2847        }
2848
2849        let new_qty = Quantity::from("0.5");
2850        let (handler, events) = subscribe_order_topic(strategy_id);
2851
2852        algo.modify_order_in_place(&mut order, Some(new_qty), None, None)
2853            .unwrap();
2854
2855        msgbus::unsubscribe_order_events(format!("events.order.{strategy_id}").into(), &handler);
2856        let events = events.borrow();
2857
2858        assert_eq!(order.quantity(), new_qty);
2859        assert_eq!(events.len(), 1);
2860        assert!(matches!(
2861            &events[0],
2862            OrderEventAny::Updated(event) if event.quantity == new_qty
2863        ));
2864    }
2865
2866    #[rstest]
2867    fn test_algorithm_modify_order_in_place_rejects_no_changes() {
2868        let mut algo = create_test_algorithm();
2869        register_algorithm(&mut algo);
2870
2871        let mut order = OrderAny::Limit(LimitOrder::new(
2872            TraderId::from("TRADER-001"),
2873            StrategyId::from("STRAT-001"),
2874            InstrumentId::from("BTC/USDT.BINANCE"),
2875            ClientOrderId::from("O-001"),
2876            OrderSide::Buy,
2877            Quantity::from("1.0"),
2878            Price::from("50000.0"),
2879            TimeInForce::Gtc,
2880            None,
2881            false,
2882            false,
2883            false,
2884            None,
2885            None,
2886            None,
2887            None,
2888            None,
2889            None,
2890            None,
2891            None,
2892            None,
2893            None,
2894            None,
2895            UUID4::new(),
2896            0.into(),
2897        ));
2898
2899        // Try to modify with same quantity - should fail
2900        let result =
2901            algo.modify_order_in_place(&mut order, Some(Quantity::from("1.0")), None, None);
2902
2903        assert!(result.is_err());
2904        assert!(
2905            result
2906                .unwrap_err()
2907                .to_string()
2908                .contains("no parameters differ")
2909        );
2910    }
2911
2912    #[rstest]
2913    fn test_spawned_order_denied_restores_primary_quantity() {
2914        let mut algo = create_test_algorithm();
2915        register_algorithm(&mut algo);
2916
2917        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
2918        let exec_algorithm_id = algo.id();
2919        let client_order_id = ClientOrderId::from("O-001");
2920
2921        let mut primary = OrderAny::Market(MarketOrder::new(
2922            TraderId::from("TRADER-001"),
2923            StrategyId::from("STRAT-001"),
2924            instrument_id,
2925            client_order_id,
2926            OrderSide::Buy,
2927            Quantity::from("1.0"),
2928            TimeInForce::Gtc,
2929            UUID4::new(),
2930            0.into(),
2931            false,
2932            false,
2933            None,
2934            None,
2935            None,
2936            None,
2937            Some(exec_algorithm_id),
2938            None,
2939            Some(client_order_id),
2940            None,
2941        ));
2942
2943        {
2944            let cache_rc = algo.core.cache_rc();
2945            let mut cache = cache_rc.borrow_mut();
2946            cache.add_order(primary.clone(), None, None, false).unwrap();
2947        }
2948
2949        let spawned = algo.spawn_market(
2950            &mut primary,
2951            Quantity::from("0.5"),
2952            TimeInForce::Fok,
2953            false,
2954            None,
2955            true,
2956        );
2957
2958        assert_eq!(primary.quantity(), Quantity::from("0.5"));
2959
2960        let spawned_order = OrderAny::Market(spawned);
2961        {
2962            let cache_rc = algo.core.cache_rc();
2963            let mut cache = cache_rc.borrow_mut();
2964            cache
2965                .add_order(spawned_order.clone(), None, None, false)
2966                .unwrap();
2967        }
2968
2969        let denied = OrderDeniedSpec::builder()
2970            .trader_id(spawned_order.trader_id())
2971            .strategy_id(spawned_order.strategy_id())
2972            .instrument_id(spawned_order.instrument_id())
2973            .client_order_id(spawned_order.client_order_id())
2974            .reason("TEST_DENIAL".into())
2975            .build();
2976
2977        {
2978            let cache_rc = algo.core.cache_rc();
2979            let mut cache = cache_rc.borrow_mut();
2980            cache.update_order(&OrderEventAny::Denied(denied)).unwrap();
2981        }
2982
2983        algo.handle_order_event(OrderEventAny::Denied(denied));
2984
2985        let restored_primary = algo.cache().order(&client_order_id).unwrap();
2986        assert_eq!(restored_primary.quantity(), Quantity::from("1.0"));
2987    }
2988
2989    #[rstest]
2990    fn test_spawned_order_rejected_restores_primary_quantity() {
2991        let mut algo = create_test_algorithm();
2992        register_algorithm(&mut algo);
2993
2994        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
2995        let exec_algorithm_id = algo.id();
2996        let client_order_id = ClientOrderId::from("O-001");
2997
2998        let mut primary = OrderAny::Market(MarketOrder::new(
2999            TraderId::from("TRADER-001"),
3000            StrategyId::from("STRAT-001"),
3001            instrument_id,
3002            client_order_id,
3003            OrderSide::Buy,
3004            Quantity::from("1.0"),
3005            TimeInForce::Gtc,
3006            UUID4::new(),
3007            0.into(),
3008            false,
3009            false,
3010            None,
3011            None,
3012            None,
3013            None,
3014            Some(exec_algorithm_id),
3015            None,
3016            Some(client_order_id),
3017            None,
3018        ));
3019
3020        {
3021            let cache_rc = algo.core.cache_rc();
3022            let mut cache = cache_rc.borrow_mut();
3023            cache.add_order(primary.clone(), None, None, false).unwrap();
3024        }
3025
3026        let spawned = algo.spawn_market(
3027            &mut primary,
3028            Quantity::from("0.5"),
3029            TimeInForce::Fok,
3030            false,
3031            None,
3032            true,
3033        );
3034
3035        assert_eq!(primary.quantity(), Quantity::from("0.5"));
3036
3037        let spawned_order = OrderAny::Market(spawned);
3038        {
3039            let cache_rc = algo.core.cache_rc();
3040            let mut cache = cache_rc.borrow_mut();
3041            cache
3042                .add_order(spawned_order.clone(), None, None, false)
3043                .unwrap();
3044        }
3045
3046        let rejected = OrderRejectedSpec::builder()
3047            .trader_id(spawned_order.trader_id())
3048            .strategy_id(spawned_order.strategy_id())
3049            .instrument_id(spawned_order.instrument_id())
3050            .client_order_id(spawned_order.client_order_id())
3051            .account_id(AccountId::from("BINANCE-001"))
3052            .reason("TEST_REJECTION".into())
3053            .build();
3054
3055        {
3056            let cache_rc = algo.core.cache_rc();
3057            let mut cache = cache_rc.borrow_mut();
3058            cache
3059                .update_order(&OrderEventAny::Rejected(rejected))
3060                .unwrap();
3061        }
3062
3063        algo.handle_order_event(OrderEventAny::Rejected(rejected));
3064
3065        let restored_primary = algo.cache().order(&client_order_id).unwrap();
3066        assert_eq!(restored_primary.quantity(), Quantity::from("1.0"));
3067    }
3068
3069    #[rstest]
3070    fn test_spawned_order_with_reduce_primary_false_does_not_restore() {
3071        let mut algo = create_test_algorithm();
3072        register_algorithm(&mut algo);
3073
3074        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
3075        let exec_algorithm_id = algo.id();
3076        let client_order_id = ClientOrderId::from("O-001");
3077
3078        let mut primary = OrderAny::Market(MarketOrder::new(
3079            TraderId::from("TRADER-001"),
3080            StrategyId::from("STRAT-001"),
3081            instrument_id,
3082            client_order_id,
3083            OrderSide::Buy,
3084            Quantity::from("1.0"),
3085            TimeInForce::Gtc,
3086            UUID4::new(),
3087            0.into(),
3088            false,
3089            false,
3090            None,
3091            None,
3092            None,
3093            None,
3094            Some(exec_algorithm_id),
3095            None,
3096            Some(client_order_id),
3097            None,
3098        ));
3099
3100        {
3101            let cache_rc = algo.core.cache_rc();
3102            let mut cache = cache_rc.borrow_mut();
3103            cache.add_order(primary.clone(), None, None, false).unwrap();
3104        }
3105
3106        let spawned = algo.spawn_market(
3107            &mut primary,
3108            Quantity::from("0.5"),
3109            TimeInForce::Fok,
3110            false,
3111            None,
3112            false,
3113        );
3114
3115        assert_eq!(primary.quantity(), Quantity::from("1.0"));
3116
3117        let spawned_order = OrderAny::Market(spawned);
3118        {
3119            let cache_rc = algo.core.cache_rc();
3120            let mut cache = cache_rc.borrow_mut();
3121            cache
3122                .add_order(spawned_order.clone(), None, None, false)
3123                .unwrap();
3124        }
3125
3126        let denied = OrderDeniedSpec::builder()
3127            .trader_id(spawned_order.trader_id())
3128            .strategy_id(spawned_order.strategy_id())
3129            .instrument_id(spawned_order.instrument_id())
3130            .client_order_id(spawned_order.client_order_id())
3131            .reason("TEST_DENIAL".into())
3132            .build();
3133
3134        {
3135            let cache_rc = algo.core.cache_rc();
3136            let mut cache = cache_rc.borrow_mut();
3137            cache.update_order(&OrderEventAny::Denied(denied)).unwrap();
3138        }
3139
3140        algo.handle_order_event(OrderEventAny::Denied(denied));
3141
3142        let final_primary = algo.cache().order(&client_order_id).unwrap();
3143        assert_eq!(final_primary.quantity(), Quantity::from("1.0"));
3144    }
3145
3146    #[rstest]
3147    fn test_multiple_spawns_with_one_denied_restores_correctly() {
3148        let mut algo = create_test_algorithm();
3149        register_algorithm(&mut algo);
3150
3151        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
3152        let exec_algorithm_id = algo.id();
3153        let client_order_id = ClientOrderId::from("O-001");
3154
3155        let mut primary = OrderAny::Market(MarketOrder::new(
3156            TraderId::from("TRADER-001"),
3157            StrategyId::from("STRAT-001"),
3158            instrument_id,
3159            client_order_id,
3160            OrderSide::Buy,
3161            Quantity::from("1.0"),
3162            TimeInForce::Gtc,
3163            UUID4::new(),
3164            0.into(),
3165            false,
3166            false,
3167            None,
3168            None,
3169            None,
3170            None,
3171            Some(exec_algorithm_id),
3172            None,
3173            Some(client_order_id),
3174            None,
3175        ));
3176
3177        {
3178            let cache_rc = algo.core.cache_rc();
3179            let mut cache = cache_rc.borrow_mut();
3180            cache.add_order(primary.clone(), None, None, false).unwrap();
3181        }
3182
3183        let spawned1 = algo.spawn_market(
3184            &mut primary,
3185            Quantity::from("0.3"),
3186            TimeInForce::Fok,
3187            false,
3188            None,
3189            true,
3190        );
3191        let spawned2 = algo.spawn_market(
3192            &mut primary,
3193            Quantity::from("0.4"),
3194            TimeInForce::Fok,
3195            false,
3196            None,
3197            true,
3198        );
3199        assert_eq!(primary.quantity(), Quantity::from("0.3"));
3200
3201        let spawned_order1 = OrderAny::Market(spawned1);
3202        let spawned_order2 = OrderAny::Market(spawned2);
3203        {
3204            let cache_rc = algo.core.cache_rc();
3205            let mut cache = cache_rc.borrow_mut();
3206            cache.add_order(spawned_order1, None, None, false).unwrap();
3207            cache
3208                .add_order(spawned_order2.clone(), None, None, false)
3209                .unwrap();
3210        }
3211
3212        let denied = OrderDeniedSpec::builder()
3213            .trader_id(spawned_order2.trader_id())
3214            .strategy_id(spawned_order2.strategy_id())
3215            .instrument_id(spawned_order2.instrument_id())
3216            .client_order_id(spawned_order2.client_order_id())
3217            .reason("TEST_DENIAL".into())
3218            .build();
3219
3220        {
3221            let cache_rc = algo.core.cache_rc();
3222            let mut cache = cache_rc.borrow_mut();
3223            cache.update_order(&OrderEventAny::Denied(denied)).unwrap();
3224        }
3225
3226        let (handler, events) = subscribe_order_topic(spawned_order2.strategy_id());
3227
3228        algo.handle_order_event(OrderEventAny::Denied(denied));
3229
3230        msgbus::unsubscribe_order_events(
3231            format!("events.order.{}", spawned_order2.strategy_id()).into(),
3232            &handler,
3233        );
3234        let events = events.borrow();
3235
3236        let restored_primary = algo.cache().order(&client_order_id).unwrap();
3237        assert_eq!(restored_primary.quantity(), Quantity::from("0.7"));
3238        assert_eq!(events.len(), 1);
3239        assert!(matches!(
3240            &events[0],
3241            OrderEventAny::Updated(event) if event.quantity == Quantity::from("0.7")
3242        ));
3243    }
3244
3245    #[rstest]
3246    fn test_spawned_order_accepted_prevents_restoration() {
3247        let mut algo = create_test_algorithm();
3248        register_algorithm(&mut algo);
3249
3250        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
3251        let exec_algorithm_id = algo.id();
3252        let client_order_id = ClientOrderId::from("O-001");
3253
3254        let mut primary = OrderAny::Market(MarketOrder::new(
3255            TraderId::from("TRADER-001"),
3256            StrategyId::from("STRAT-001"),
3257            instrument_id,
3258            client_order_id,
3259            OrderSide::Buy,
3260            Quantity::from("1.0"),
3261            TimeInForce::Gtc,
3262            UUID4::new(),
3263            0.into(),
3264            false,
3265            false,
3266            None,
3267            None,
3268            None,
3269            None,
3270            Some(exec_algorithm_id),
3271            None,
3272            Some(client_order_id),
3273            None,
3274        ));
3275
3276        {
3277            let cache_rc = algo.core.cache_rc();
3278            let mut cache = cache_rc.borrow_mut();
3279            cache.add_order(primary.clone(), None, None, false).unwrap();
3280        }
3281
3282        let spawned = algo.spawn_market(
3283            &mut primary,
3284            Quantity::from("0.5"),
3285            TimeInForce::Fok,
3286            false,
3287            None,
3288            true,
3289        );
3290
3291        assert_eq!(primary.quantity(), Quantity::from("0.5"));
3292
3293        let mut spawned_order = OrderAny::Market(spawned);
3294        {
3295            let cache_rc = algo.core.cache_rc();
3296            let mut cache = cache_rc.borrow_mut();
3297            cache
3298                .add_order(spawned_order.clone(), None, None, false)
3299                .unwrap();
3300        }
3301
3302        let accepted = OrderAcceptedSpec::builder()
3303            .trader_id(spawned_order.trader_id())
3304            .strategy_id(spawned_order.strategy_id())
3305            .instrument_id(spawned_order.instrument_id())
3306            .client_order_id(spawned_order.client_order_id())
3307            .venue_order_id(VenueOrderId::from("V-123"))
3308            .account_id(AccountId::from("BINANCE-001"))
3309            .build();
3310
3311        {
3312            let cache_rc = algo.core.cache_rc();
3313            let mut cache = cache_rc.borrow_mut();
3314            spawned_order = cache
3315                .update_order(&OrderEventAny::Accepted(accepted))
3316                .unwrap();
3317        }
3318
3319        algo.handle_order_event(OrderEventAny::Accepted(accepted));
3320
3321        let primary_after_accept = algo.cache().order(&client_order_id).unwrap();
3322        assert_eq!(primary_after_accept.quantity(), Quantity::from("0.5"));
3323
3324        // Cancel after acceptance - no restoration should occur
3325        let canceled = OrderCanceledSpec::builder()
3326            .trader_id(spawned_order.trader_id())
3327            .strategy_id(spawned_order.strategy_id())
3328            .instrument_id(spawned_order.instrument_id())
3329            .client_order_id(spawned_order.client_order_id())
3330            .venue_order_id(VenueOrderId::from("V-123"))
3331            .account_id(AccountId::from("BINANCE-001"))
3332            .build();
3333
3334        {
3335            let cache_rc = algo.core.cache_rc();
3336            let mut cache = cache_rc.borrow_mut();
3337            cache
3338                .update_order(&OrderEventAny::Canceled(canceled))
3339                .unwrap();
3340        }
3341
3342        algo.handle_order_event(OrderEventAny::Canceled(canceled));
3343
3344        let final_primary = algo.cache().order(&client_order_id).unwrap();
3345        assert_eq!(final_primary.quantity(), Quantity::from("0.5"));
3346    }
3347
3348    #[rstest]
3349    #[should_panic(expected = "exceeds primary leaves_qty")]
3350    fn test_spawn_quantity_exceeds_leaves_qty_panics() {
3351        let mut algo = create_test_algorithm();
3352        register_algorithm(&mut algo);
3353
3354        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
3355        let exec_algorithm_id = algo.id();
3356        let client_order_id = ClientOrderId::from("O-001");
3357
3358        let mut primary = OrderAny::Market(MarketOrder::new(
3359            TraderId::from("TRADER-001"),
3360            StrategyId::from("STRAT-001"),
3361            instrument_id,
3362            client_order_id,
3363            OrderSide::Buy,
3364            Quantity::from("1.0"),
3365            TimeInForce::Gtc,
3366            UUID4::new(),
3367            0.into(),
3368            false,
3369            false,
3370            None,
3371            None,
3372            None,
3373            None,
3374            Some(exec_algorithm_id),
3375            None,
3376            Some(client_order_id),
3377            None,
3378        ));
3379
3380        {
3381            let cache_rc = algo.core.cache_rc();
3382            let mut cache = cache_rc.borrow_mut();
3383            cache.add_order(primary.clone(), None, None, false).unwrap();
3384        }
3385
3386        let _ = algo.spawn_market(
3387            &mut primary,
3388            Quantity::from("0.8"),
3389            TimeInForce::Fok,
3390            false,
3391            None,
3392            true,
3393        );
3394
3395        assert_eq!(primary.quantity(), Quantity::from("0.2"));
3396        assert_eq!(primary.leaves_qty(), Quantity::from("0.2"));
3397
3398        // Should panic - spawning 0.5 when only 0.2 leaves_qty remains
3399        let _ = algo.spawn_market(
3400            &mut primary,
3401            Quantity::from("0.5"),
3402            TimeInForce::Fok,
3403            false,
3404            None,
3405            true,
3406        );
3407    }
3408}