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