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