1use 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
74pub trait ExecutionAlgorithm: DataActor {
101 fn id(&self) -> ExecAlgorithmId
103 where
104 Self: ExecutionAlgorithmNative,
105 {
106 ExecutionAlgorithmNative::exec_algorithm_core(self).exec_algorithm_id
107 }
108
109 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 fn on_order(&mut self, order: OrderAny) -> anyhow::Result<()>;
168
169 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 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 ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
250 .remove_submit_params(&order.client_order_id());
251
252 Ok(())
253 }
254
255 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 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 log::warn!(
361 "Cannot modify {}: order is not active-local",
362 command.client_order_id
363 );
364 Ok(())
365 }
366
367 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, order.venue_order_id(),
386 order.account_id(),
387 )
388 }
389
390 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, order.venue_order_id(),
410 )
411 }
412
413 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, order.venue_order_id(),
433 )
434 }
435
436 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 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 #[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 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, 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 #[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 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, 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 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, primary.venue_order_id(),
674 primary.account_id(),
675 None, None, None, 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 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 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, primary.venue_order_id(),
756 primary.account_id(),
757 None, None, None, 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 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 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, );
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 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, None, );
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 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 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 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 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, order.venue_order_id(),
1059 order.account_id(),
1060 price,
1061 trigger_price,
1062 None, 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 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, None, );
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 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 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 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 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 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 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 fn on_stop(&mut self) -> anyhow::Result<()> {
1379 Ok(())
1380 }
1381
1382 fn on_resume(&mut self) -> anyhow::Result<()> {
1388 Ok(())
1389 }
1390
1391 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 fn on_time_event(&mut self, _event: &TimeEvent) -> anyhow::Result<()> {
1413 Ok(())
1414 }
1415
1416 #[allow(unused_variables)]
1418 fn on_order_initialized(&mut self, event: OrderInitialized) {}
1419
1420 #[allow(unused_variables)]
1422 fn on_order_denied(&mut self, event: OrderDenied) {}
1423
1424 #[allow(unused_variables)]
1426 fn on_order_emulated(&mut self, event: OrderEmulated) {}
1427
1428 #[allow(unused_variables)]
1430 fn on_order_released(&mut self, event: OrderReleased) {}
1431
1432 #[allow(unused_variables)]
1434 fn on_order_submitted(&mut self, event: OrderSubmitted) {}
1435
1436 #[allow(unused_variables)]
1438 fn on_order_rejected(&mut self, event: OrderRejected) {}
1439
1440 #[allow(unused_variables)]
1442 fn on_order_accepted(&mut self, event: OrderAccepted) {}
1443
1444 #[allow(unused_variables)]
1446 fn on_algo_order_canceled(&mut self, event: OrderCanceled) {}
1447
1448 #[allow(unused_variables)]
1450 fn on_order_expired(&mut self, event: OrderExpired) {}
1451
1452 #[allow(unused_variables)]
1454 fn on_order_triggered(&mut self, event: OrderTriggered) {}
1455
1456 #[allow(unused_variables)]
1458 fn on_order_pending_update(&mut self, event: OrderPendingUpdate) {}
1459
1460 #[allow(unused_variables)]
1462 fn on_order_pending_cancel(&mut self, event: OrderPendingCancel) {}
1463
1464 #[allow(unused_variables)]
1466 fn on_order_modify_rejected(&mut self, event: OrderModifyRejected) {}
1467
1468 #[allow(unused_variables)]
1470 fn on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {}
1471
1472 #[allow(unused_variables)]
1474 fn on_order_updated(&mut self, event: OrderUpdated) {}
1475
1476 #[allow(unused_variables)]
1478 fn on_algo_order_filled(&mut self, event: OrderFilled) {}
1479
1480 #[allow(unused_variables)]
1482 fn on_order_fill_voided(&mut self, event: &OrderFillVoided) {}
1483
1484 #[allow(unused_variables)]
1486 fn on_order_event(&mut self, event: OrderEventAny) {}
1487
1488 #[allow(unused_variables)]
1490 fn on_position_opened(&mut self, event: PositionOpened) {}
1491
1492 #[allow(unused_variables)]
1494 fn on_position_changed(&mut self, event: PositionChanged) {}
1495
1496 #[allow(unused_variables)]
1498 fn on_position_closed(&mut self, event: PositionClosed) {}
1499
1500 #[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 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 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 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, false, None, None, None, None, None, None, None, None, ));
2037
2038 let spawned = algo.spawn_market(
2039 &mut primary,
2040 Quantity::from("0.5"),
2041 TimeInForce::Ioc,
2042 false,
2043 None, false, );
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 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, false, false, None, None, None, false, );
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, false, None, None, None, false, );
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, true, None, None, Some(linked_order_ids.clone()),
2338 None, 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, 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, false, false, None, None, None, 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, false, None, None, None, 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 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 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, );
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, );
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, false, false, false, None, None, None, None, None, None, None, None, None, None, None, 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 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 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 let _ = algo.spawn_market(
3798 &mut primary,
3799 Quantity::from("0.5"),
3800 TimeInForce::Fok,
3801 false,
3802 None,
3803 true,
3804 );
3805 }
3806}