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