1pub mod api;
17#[doc(hidden)]
18pub mod binding;
19pub mod config;
20pub mod core;
21
22pub use core::{StrategyCore, StrategyNative};
23use std::panic::{AssertUnwindSafe, catch_unwind};
24
25use ahash::AHashSet;
26pub use api::{OrderApi, PortfolioApi};
27use binding::StrategyBinding;
28pub use config::{ImportableStrategyConfig, StrategyConfig};
29use nautilus_common::{
30 actor::DataActor,
31 component::Component,
32 enums::ComponentState,
33 logging::{CMD, EVT, RECV, SEND},
34 messages::execution::{
35 BatchCancelOrders, BatchModifyOrders, CancelAllOrders, CancelOrder, ModifyOrder,
36 QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList, TradingCommand,
37 },
38 msgbus::{self, MessagingSwitchboard},
39 timer::TimeEvent,
40};
41use nautilus_core::{DurationNanos, Params, UUID4};
42use nautilus_execution::order_manager::OrderManagerAction;
43use nautilus_model::{
44 enums::{OrderSide, OrderStatus, PositionSide, TimeInForce},
45 events::{
46 OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied, OrderEmulated,
47 OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled, OrderInitialized,
48 OrderModifyRejected, OrderPendingCancel, OrderPendingUpdate, OrderRejected, OrderReleased,
49 OrderSubmitted, OrderTriggered, OrderUpdated, PositionChanged, PositionClosed,
50 PositionEvent, PositionOpened,
51 },
52 identifiers::{
53 AccountId, ClientId, ClientOrderId, ExecAlgorithmId, InstrumentId, PositionId, StrategyId,
54 TraderId,
55 },
56 orders::{
57 LIMIT_ORDER_TYPES, Order, OrderAny, OrderCore, OrderError, OrderList, STOP_ORDER_TYPES,
58 },
59 position::Position,
60 types::{Price, Quantity},
61};
62use ustr::Ustr;
63
64pub type BatchModifyOrder = (
66 ClientOrderId,
67 Option<Quantity>,
68 Option<Price>,
69 Option<Price>,
70);
71
72pub trait Strategy: DataActor {
108 fn external_order_instrument_ids(&self) -> Option<Vec<InstrumentId>> {
112 None
113 }
114
115 fn set_external_order_instrument_ids(
126 &mut self,
127 instrument_ids: Vec<InstrumentId>,
128 ) -> anyhow::Result<()>
129 where
130 Self: StrategyNative,
131 {
132 let core = StrategyNative::strategy_core_mut(self);
133 let strategy_id = registered_strategy_id(core)?;
134 if !core.actor.is_registered() {
135 anyhow::bail!("Strategy {strategy_id} is not registered with a trader");
136 }
137 let cache = core.cache_rc();
138 cache
139 .try_borrow_mut()
140 .map_err(|e| anyhow::anyhow!("Cannot set external order claims: {e}"))?
141 .set_external_order_claims(strategy_id, &instrument_ids)?;
142 core.config.external_order_instrument_ids = Some(instrument_ids);
143 Ok(())
144 }
145
146 fn strategy_id(&self) -> Option<StrategyId>
148 where
149 Self: StrategyNative,
150 {
151 StrategyNative::strategy_core(self).strategy_id()
152 }
153
154 fn order(&self) -> OrderApi<'_>
156 where
157 Self: StrategyNative,
158 {
159 StrategyNative::strategy_core(self).order()
160 }
161
162 fn portfolio(&self) -> PortfolioApi<'_>
164 where
165 Self: StrategyNative,
166 {
167 StrategyNative::strategy_core(self).portfolio_api()
168 }
169
170 fn submit_order(
176 &mut self,
177 order: OrderAny,
178 position_id: Option<PositionId>,
179 client_id: Option<ClientId>,
180 params: Option<Params>,
181 ) -> anyhow::Result<()>
182 where
183 Self: StrategyBinding,
184 {
185 self.binding_submit_order(order, position_id, client_id, params)
186 }
187
188 fn submit_order_list(
195 &mut self,
196 mut orders: Vec<OrderAny>,
197 position_id: Option<PositionId>,
198 client_id: Option<ClientId>,
199 params: Option<Params>,
200 ) -> anyhow::Result<()>
201 where
202 Self: StrategyNative,
203 {
204 if orders.is_empty() {
205 log::error!("OrderList denied: no orders to submit");
206 anyhow::bail!("OrderList denied: no orders to submit");
207 }
208
209 for order in &orders {
210 if order.status() != OrderStatus::Initialized {
211 anyhow::bail!(
212 "Order in list denied: invalid status for {}, expected INITIALIZED",
213 order.client_order_id()
214 );
215 }
216 }
217
218 let first_venue = orders[0].instrument_id().venue;
219 for order in &orders {
220 if order.instrument_id().venue != first_venue {
221 anyhow::bail!(
222 "OrderList denied: orders must share the same venue; \
223 expected {first_venue}, found {} on {}",
224 order.instrument_id().venue,
225 order.client_order_id(),
226 );
227 }
228 }
229
230 let should_deny = {
231 let core = StrategyNative::strategy_core_mut(self);
232 let tag = core.market_exit_tag;
233 core.is_exiting
234 && orders.iter().any(|o| {
235 !o.is_reduce_only() && !o.tags().is_some_and(|tags| tags.contains(&tag))
236 })
237 };
238
239 if should_deny {
240 self.deny_order_list(&orders, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
241 return Ok(());
242 }
243
244 let core = StrategyNative::strategy_core_mut(self);
245
246 let trader_id = registered_trader_id(core)?;
247 let strategy_id = registered_strategy_id(core)?;
248 let ts_init = core.clock_mut().timestamp_ns();
249
250 let order_list = if orders.first().is_some_and(|o| o.order_list_id().is_some()) {
252 OrderList::from_orders(&orders, ts_init)
253 } else {
254 core.order_factory().create_list(&mut orders, ts_init)
255 };
256
257 if let Err(e) = order_list.validate() {
258 log::error!("OrderList denied: {e}");
259 anyhow::bail!("OrderList denied: {e}");
260 }
261
262 {
263 let cache_rc = core.cache_rc();
264 let mut cache = cache_rc.try_borrow_mut().map_err(|_| {
265 anyhow::anyhow!(
266 "Cannot submit order list {}: cache is currently borrowed",
267 order_list.id
268 )
269 })?;
270
271 if cache.order_list_exists(&order_list.id) {
272 anyhow::bail!("OrderList denied: duplicate {}", order_list.id);
273 }
274
275 for order in &orders {
276 if cache.order_exists(&order.client_order_id()) {
277 anyhow::bail!(
278 "Order in list denied: duplicate {}",
279 order.client_order_id()
280 );
281 }
282 }
283
284 cache.add_order_list(order_list.clone())?;
285 for order in &orders {
286 cache.add_order(order.clone(), position_id, client_id, true)?;
287 }
288 }
289
290 for order in &orders {
291 publish_order_initialized(order);
292 }
293
294 let params = params.filter(|params| !params.is_empty());
295
296 let first_order = orders.first();
297 let order_inits: Vec<_> = orders.iter().map(|o| o.init_event().clone()).collect();
298 let exec_algorithm_id = first_order.and_then(Order::exec_algorithm_id);
299
300 let command = SubmitOrderList::new(
301 trader_id,
302 client_id,
303 strategy_id,
304 order_list,
305 order_inits,
306 exec_algorithm_id,
307 position_id,
308 params,
309 UUID4::new(),
310 ts_init,
311 None, );
313
314 let has_emulated_order = orders
315 .iter()
316 .any(|o| o.emulation_trigger().is_some() || o.is_emulated());
317
318 if has_emulated_order {
319 send_emulator_command(TradingCommand::SubmitOrderList(command));
320 } else if let Some(algo_id) = exec_algorithm_id {
321 let endpoint = format!("{algo_id}.execute");
322 msgbus::send_any(endpoint.into(), &TradingCommand::SubmitOrderList(command));
323 } else {
324 send_risk_command(TradingCommand::SubmitOrderList(command));
325 }
326
327 for order in &orders {
328 self.set_gtd_expiry(order)?;
329 }
330
331 Ok(())
332 }
333
334 fn modify_order(
340 &mut self,
341 client_order_id: ClientOrderId,
342 quantity: Option<Quantity>,
343 price: Option<Price>,
344 trigger_price: Option<Price>,
345 client_id: Option<ClientId>,
346 params: Option<Params>,
347 ) -> anyhow::Result<()>
348 where
349 Self: StrategyNative,
350 {
351 let (trader_id, strategy_id) = {
352 let core = StrategyNative::strategy_core_mut(self);
353 (registered_trader_id(core)?, registered_strategy_id(core)?)
354 };
355
356 let params = params.filter(|params| !params.is_empty());
357
358 let order = StrategyNative::strategy_core_mut(self)
360 .cache_rc()
361 .borrow()
362 .try_order_owned(&client_order_id)
363 .map_err(|e| anyhow::anyhow!("Cannot modify order: {e}"))?;
364
365 let mut updating = false;
366
367 if quantity.is_some_and(|q| q != order.quantity() || order.is_pending_update()) {
368 updating = true;
369 }
370
371 if let Some(price) = price {
372 if !LIMIT_ORDER_TYPES.contains(&order.order_type()) {
373 anyhow::bail!("{} orders do not have a LIMIT price", order.order_type());
374 }
375
376 if Some(price) != order.price() {
377 updating = true;
378 }
379 }
380
381 if let Some(trigger_price) = trigger_price {
382 if !STOP_ORDER_TYPES.contains(&order.order_type()) {
383 anyhow::bail!(
384 "{} orders do not have a STOP trigger price",
385 order.order_type()
386 );
387 }
388
389 if Some(trigger_price) != order.trigger_price() {
390 updating = true;
391 }
392 }
393
394 if !updating {
395 log::error!(
396 "Cannot create command ModifyOrder: quantity, price, and trigger were either None \
397 or the same as existing values"
398 );
399 return Ok(());
400 }
401
402 if order.is_closed() || order.is_pending_cancel() {
403 log::warn!(
404 "Cannot create command ModifyOrder: state is {:?}, {order:?}",
405 order.status()
406 );
407 return Ok(());
408 }
409
410 if !self.mark_order_pending_update(&order)? {
411 return Ok(());
412 }
413
414 let command = ModifyOrder::new(
415 trader_id,
416 client_id,
417 strategy_id,
418 order.instrument_id(),
419 order.client_order_id(),
420 order.venue_order_id(),
421 quantity,
422 price,
423 trigger_price,
424 UUID4::new(),
425 StrategyNative::strategy_core_mut(self)
426 .clock_mut()
427 .timestamp_ns(),
428 params,
429 None, );
431
432 if order.emulation_trigger().is_some() || order.is_emulated() {
433 send_emulator_command(TradingCommand::ModifyOrder(command));
434 } else if let Some(algo_id) = order
435 .exec_algorithm_id()
436 .filter(|_| order.is_active_local())
437 {
438 let endpoint = format!("{algo_id}.execute");
439 msgbus::send_any(endpoint.into(), &TradingCommand::ModifyOrder(command));
440 } else {
441 send_risk_command(TradingCommand::ModifyOrder(command));
442 }
443 Ok(())
444 }
445
446 fn modify_orders(
455 &mut self,
456 updates: Vec<BatchModifyOrder>,
457 client_id: Option<ClientId>,
458 params: Option<Params>,
459 ) -> anyhow::Result<()>
460 where
461 Self: StrategyNative,
462 {
463 if updates.is_empty() {
464 anyhow::bail!("Cannot batch modify empty order list");
465 }
466
467 let (trader_id, strategy_id, ts_init) = {
468 let core = StrategyNative::strategy_core_mut(self);
469 (
470 registered_trader_id(core)?,
471 registered_strategy_id(core)?,
472 core.clock_mut().timestamp_ns(),
473 )
474 };
475
476 let orders: Vec<OrderAny> = {
477 let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
478 let cache = cache_rc.borrow();
479 updates
480 .iter()
481 .map(|(client_order_id, _, _, _)| {
482 cache
483 .try_order_owned(client_order_id)
484 .map_err(|e| anyhow::anyhow!("Cannot modify order: {e}"))
485 })
486 .collect::<Result<_, _>>()?
487 };
488
489 let instrument_id = orders[0].instrument_id();
490
491 for (order, (_, quantity, price, trigger_price)) in orders.iter().zip(updates.iter()) {
492 if order.instrument_id() != instrument_id {
493 anyhow::bail!(
494 "Cannot batch modify orders for different instruments: {} vs {}",
495 instrument_id,
496 order.instrument_id()
497 );
498 }
499
500 if order.is_emulated() || order.is_active_local() {
501 anyhow::bail!("Cannot include emulated or local orders in batch modify");
502 }
503
504 let mut updating = false;
505
506 if quantity.is_some_and(|q| q != order.quantity()) {
507 updating = true;
508 }
509
510 if let Some(price) = price {
511 if !LIMIT_ORDER_TYPES.contains(&order.order_type()) {
512 anyhow::bail!("{} orders do not have a LIMIT price", order.order_type());
513 }
514
515 if Some(*price) != order.price() {
516 updating = true;
517 }
518 }
519
520 if let Some(trigger_price) = trigger_price {
521 if !STOP_ORDER_TYPES.contains(&order.order_type()) {
522 anyhow::bail!(
523 "{} orders do not have a STOP trigger price",
524 order.order_type()
525 );
526 }
527
528 if Some(*trigger_price) != order.trigger_price() {
529 updating = true;
530 }
531 }
532
533 if !updating {
534 anyhow::bail!(
535 "Cannot create command BatchModifyOrders: quantity, price, and trigger were \
536 either None or the same as existing values for {}",
537 order.client_order_id()
538 );
539 }
540
541 if order.is_closed() || order.is_pending_cancel() {
542 anyhow::bail!(
543 "Cannot create command BatchModifyOrders: state is {:?}, {order:?}",
544 order.status()
545 );
546 }
547 }
548
549 let params = params.filter(|params| !params.is_empty());
550 let mut modifies = Vec::with_capacity(orders.len());
551
552 for (order, (_, quantity, price, trigger_price)) in orders.into_iter().zip(updates) {
553 if !self.mark_order_pending_update(&order)? {
554 continue;
555 }
556
557 modifies.push(ModifyOrder::new(
558 trader_id,
559 client_id,
560 strategy_id,
561 instrument_id,
562 order.client_order_id(),
563 order.venue_order_id(),
564 quantity,
565 price,
566 trigger_price,
567 UUID4::new(),
568 ts_init,
569 params.clone(),
570 None, ));
572 }
573
574 if modifies.is_empty() {
575 log::warn!("Cannot send `BatchModifyOrders`, no valid modify commands");
576 return Ok(());
577 }
578
579 let command = BatchModifyOrders::new(
580 trader_id,
581 client_id,
582 strategy_id,
583 instrument_id,
584 modifies,
585 UUID4::new(),
586 ts_init,
587 params,
588 None, );
590
591 send_risk_command(TradingCommand::ModifyOrders(command));
592 Ok(())
593 }
594
595 fn cancel_order(
601 &mut self,
602 client_order_id: ClientOrderId,
603 client_id: Option<ClientId>,
604 params: Option<Params>,
605 ) -> anyhow::Result<()>
606 where
607 Self: StrategyNative,
608 {
609 let (trader_id, strategy_id, ts_init) = {
610 let core = StrategyNative::strategy_core_mut(self);
611 (
612 registered_trader_id(core)?,
613 registered_strategy_id(core)?,
614 core.clock_mut().timestamp_ns(),
615 )
616 };
617
618 let params = params.filter(|params| !params.is_empty());
619
620 let order = StrategyNative::strategy_core_mut(self)
624 .cache_rc()
625 .borrow()
626 .try_order_owned(&client_order_id)
627 .map_err(|e| anyhow::anyhow!("Cannot cancel order: {e}"))?;
628
629 if !self.mark_order_pending_cancel(&order)? {
630 return Ok(());
631 }
632
633 let command = CancelOrder::new(
634 trader_id,
635 client_id,
636 strategy_id,
637 order.instrument_id(),
638 order.client_order_id(),
639 order.venue_order_id(),
640 UUID4::new(),
641 ts_init,
642 params,
643 None, );
645
646 if order.emulation_trigger().is_some() || order.is_emulated() {
647 send_emulator_command(TradingCommand::CancelOrder(command));
648 } else if let Some(algo_id) = order
649 .exec_algorithm_id()
650 .filter(|_| order.is_active_local())
651 {
652 let endpoint = format!("{algo_id}.execute");
653 msgbus::send_any(endpoint.into(), &TradingCommand::CancelOrder(command));
654 } else {
655 send_exec_command(TradingCommand::CancelOrder(command));
656 }
657
658 if StrategyNative::strategy_core(self).config.manage_gtd_expiry
659 && order.time_in_force() == TimeInForce::Gtd
660 && self.has_gtd_expiry_timer(&order.client_order_id())
661 {
662 self.cancel_gtd_expiry(&order.client_order_id());
663 }
664
665 Ok(())
666 }
667
668 fn cancel_orders(
675 &mut self,
676 client_order_ids: Vec<ClientOrderId>,
677 client_id: Option<ClientId>,
678 params: Option<Params>,
679 ) -> anyhow::Result<()>
680 where
681 Self: StrategyNative,
682 {
683 if client_order_ids.is_empty() {
684 anyhow::bail!("Cannot batch cancel empty order list");
685 }
686
687 let (trader_id, strategy_id, ts_init) = {
688 let core = StrategyNative::strategy_core_mut(self);
689 (
690 registered_trader_id(core)?,
691 registered_strategy_id(core)?,
692 core.clock_mut().timestamp_ns(),
693 )
694 };
695
696 let orders: Vec<OrderAny> = {
698 let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
699 let cache = cache_rc.borrow();
700 client_order_ids
701 .iter()
702 .map(|id| {
703 cache
704 .try_order_owned(id)
705 .map_err(|e| anyhow::anyhow!("Cannot cancel order: {e}"))
706 })
707 .collect::<Result<_, _>>()?
708 };
709
710 let instrument_id = orders[0].instrument_id();
711
712 for order in &orders {
713 if order.instrument_id() != instrument_id {
714 anyhow::bail!(
715 "Cannot batch cancel orders for different instruments: {} vs {}",
716 instrument_id,
717 order.instrument_id()
718 );
719 }
720
721 if order.is_emulated() || order.is_active_local() {
722 anyhow::bail!("Cannot include emulated or local orders in batch cancel");
723 }
724 }
725
726 let mut cancels = Vec::with_capacity(orders.len());
727
728 for order in orders {
729 if !self.mark_order_pending_cancel(&order)? {
730 continue;
731 }
732
733 cancels.push(CancelOrder::new(
734 trader_id,
735 client_id,
736 strategy_id,
737 instrument_id,
738 order.client_order_id(),
739 order.venue_order_id(),
740 UUID4::new(),
741 ts_init,
742 params.clone(),
743 None, ));
745 }
746
747 if cancels.is_empty() {
748 log::warn!("Cannot send `BatchCancelOrders`, no valid cancel commands");
749 return Ok(());
750 }
751
752 let command = BatchCancelOrders::new(
753 trader_id,
754 client_id,
755 strategy_id,
756 instrument_id,
757 cancels,
758 UUID4::new(),
759 ts_init,
760 params,
761 None, );
763
764 send_exec_command(TradingCommand::CancelOrders(command));
765 Ok(())
766 }
767
768 fn mark_order_pending_update(&mut self, order: &OrderAny) -> anyhow::Result<bool>
774 where
775 Self: StrategyNative,
776 {
777 if order.is_active_local() || order.is_pending_update() {
778 return Ok(true);
779 }
780
781 let strategy_id = order.strategy_id();
782 required_account_id(order, "pending update")?;
783 let event = OrderEventAny::PendingUpdate(self.generate_order_pending_update(order));
784
785 {
786 let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
787 let mut cache = cache_rc.borrow_mut();
788 match cache.update_order(&event) {
789 Ok(_) => {}
790 Err(e)
791 if matches!(
792 e.downcast_ref::<OrderError>(),
793 Some(OrderError::InvalidStateTransition)
794 ) =>
795 {
796 log::warn!("InvalidStateTrigger: {e}, did not apply pending update event");
797 return Ok(false);
798 }
799 Err(e) => return Err(e),
800 }
801 }
802
803 let topic = format!("events.order.{strategy_id}");
804 msgbus::publish_order_event(topic.into(), &event);
805 msgbus::publish_order_event(
806 msgbus::switchboard::get_order_pending_update_topic(order.instrument_id()),
807 &event,
808 );
809
810 Ok(true)
811 }
812
813 fn mark_order_pending_cancel(&mut self, order: &OrderAny) -> anyhow::Result<bool>
819 where
820 Self: StrategyNative,
821 {
822 if order.is_closed() || order.is_pending_cancel() {
823 log::warn!(
824 "Cannot cancel order: state is {:?}, {order:?}",
825 order.status()
826 );
827 return Ok(false);
828 }
829
830 if order.is_active_local() {
831 return Ok(true);
832 }
833
834 let strategy_id = order.strategy_id();
835 required_account_id(order, "pending cancel")?;
836 let event = OrderEventAny::PendingCancel(self.generate_order_pending_cancel(order));
837
838 {
839 let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
840 let mut cache = cache_rc.borrow_mut();
841 match cache.update_order(&event) {
842 Ok(_) => {}
843 Err(e)
844 if matches!(
845 e.downcast_ref::<OrderError>(),
846 Some(OrderError::InvalidStateTransition)
847 ) =>
848 {
849 log::warn!("InvalidStateTrigger: {e}, did not apply pending cancel event");
850 return Ok(false);
851 }
852 Err(e) => return Err(e),
853 }
854 cache.update_order_pending_cancel_local(order);
855 }
856
857 let topic = format!("events.order.{strategy_id}");
858 msgbus::publish_order_event(topic.into(), &event);
859 msgbus::publish_order_event(
860 msgbus::switchboard::get_order_pending_cancel_topic(order.instrument_id()),
861 &event,
862 );
863
864 Ok(true)
865 }
866
867 fn generate_order_pending_update(&mut self, order: &OrderAny) -> OrderPendingUpdate
869 where
870 Self: StrategyNative,
871 {
872 let ts_now = StrategyNative::strategy_core_mut(self)
873 .clock_mut()
874 .timestamp_ns();
875
876 OrderPendingUpdate::new(
877 order.trader_id(),
878 order.strategy_id(),
879 order.instrument_id(),
880 order.client_order_id(),
881 order.account_id(),
882 UUID4::new(),
883 ts_now,
884 ts_now,
885 false,
886 order.venue_order_id(),
887 )
888 }
889
890 fn generate_order_pending_cancel(&mut self, order: &OrderAny) -> OrderPendingCancel
892 where
893 Self: StrategyNative,
894 {
895 let ts_now = StrategyNative::strategy_core_mut(self)
896 .clock_mut()
897 .timestamp_ns();
898
899 OrderPendingCancel::new(
900 order.trader_id(),
901 order.strategy_id(),
902 order.instrument_id(),
903 order.client_order_id(),
904 order.account_id(),
905 UUID4::new(),
906 ts_now,
907 ts_now,
908 false,
909 order.venue_order_id(),
910 )
911 }
912
913 fn cancel_all_orders(
925 &mut self,
926 instrument_id: InstrumentId,
927 order_side: Option<OrderSide>,
928 client_id: Option<ClientId>,
929 strategy_only: bool,
930 params: Option<Params>,
931 ) -> anyhow::Result<()>
932 where
933 Self: StrategyNative,
934 {
935 let params = params.filter(|params| !params.is_empty());
936 let core = StrategyNative::strategy_core_mut(self);
937
938 let trader_id = registered_trader_id(core)?;
939 let strategy_id = registered_strategy_id(core)?;
940 let ts_init = core.clock_mut().timestamp_ns();
941
942 if !strategy_only {
943 let command_id = UUID4::new();
944 let command = CancelAllOrders::new(
945 trader_id,
946 client_id,
947 strategy_id,
948 instrument_id,
949 order_side,
950 command_id,
951 ts_init,
952 params,
953 Some(command_id),
954 );
955
956 send_exec_command(TradingCommand::CancelAllOrders(command));
957 return Ok(());
958 }
959
960 let cache = core.cache_ref();
961
962 let mut open_order_ids: Vec<ClientOrderId> = cache
963 .orders_open(
964 None,
965 Some(&instrument_id),
966 Some(&strategy_id),
967 None,
968 order_side,
969 )
970 .into_iter()
971 .map(|order| order.client_order_id())
972 .collect();
973
974 let mut emulated_order_ids: Vec<ClientOrderId> = cache
975 .orders_emulated(
976 None,
977 Some(&instrument_id),
978 Some(&strategy_id),
979 None,
980 order_side,
981 )
982 .into_iter()
983 .map(|order| order.client_order_id())
984 .collect();
985
986 let mut inflight_order_ids: Vec<ClientOrderId> = cache
987 .orders_inflight(
988 None,
989 Some(&instrument_id),
990 Some(&strategy_id),
991 None,
992 order_side,
993 )
994 .into_iter()
995 .map(|order| order.client_order_id())
996 .collect();
997
998 let mut exec_algorithm_ids: Vec<_> = cache.exec_algorithm_ids().into_iter().collect();
1002 exec_algorithm_ids.sort();
1003 let mut algo_order_ids: Vec<ClientOrderId> = Vec::new();
1004
1005 for algo_id in &exec_algorithm_ids {
1006 algo_order_ids.extend(
1007 cache
1008 .orders_for_exec_algorithm(
1009 algo_id,
1010 None,
1011 Some(&instrument_id),
1012 Some(&strategy_id),
1013 None,
1014 order_side,
1015 )
1016 .into_iter()
1017 .map(|order| order.client_order_id()),
1018 );
1019 }
1020
1021 let matches_client = |client_order_id: &ClientOrderId| {
1022 client_id.is_none_or(|client_id| {
1023 cache
1024 .client_id(client_order_id)
1025 .is_none_or(|order_client_id| *order_client_id == client_id)
1026 })
1027 };
1028
1029 open_order_ids.retain(&matches_client);
1030 emulated_order_ids.retain(&matches_client);
1031 inflight_order_ids.retain(&matches_client);
1032 algo_order_ids.retain(&matches_client);
1033
1034 let open_count = open_order_ids.len();
1035 let emulated_count = emulated_order_ids.len();
1036 let inflight_count = inflight_order_ids.len();
1037 let algo_count = algo_order_ids.len();
1038
1039 let mut cancel_routes: Vec<_> = open_order_ids
1040 .iter()
1041 .chain(&emulated_order_ids)
1042 .chain(&inflight_order_ids)
1043 .chain(&algo_order_ids)
1044 .map(|client_order_id| {
1045 (
1046 *client_order_id,
1047 client_id.or_else(|| cache.client_id(client_order_id).copied()),
1048 )
1049 })
1050 .collect();
1051 cancel_routes.sort_by_key(|(client_order_id, _)| *client_order_id);
1052 cancel_routes.dedup_by_key(|(client_order_id, _)| *client_order_id);
1053
1054 drop(cache);
1055
1056 if open_count == 0 && emulated_count == 0 && inflight_count == 0 && algo_count == 0 {
1057 let side_str = order_side.map(|s| format!(" {s}")).unwrap_or_default();
1058 log::info!("No {instrument_id} open, emulated, or inflight{side_str} orders to cancel");
1059 return Ok(());
1060 }
1061
1062 let side_str = order_side.map(|s| format!(" {s}")).unwrap_or_default();
1063
1064 if open_count > 0 {
1065 log::info!(
1066 "Canceling {open_count} open{side_str} {instrument_id} order{}",
1067 if open_count == 1 { "" } else { "s" }
1068 );
1069 }
1070
1071 if emulated_count > 0 {
1072 log::info!(
1073 "Canceling {emulated_count} emulated{side_str} {instrument_id} order{}",
1074 if emulated_count == 1 { "" } else { "s" }
1075 );
1076 }
1077
1078 if inflight_count > 0 {
1079 log::info!(
1080 "Canceling {inflight_count} inflight{side_str} {instrument_id} order{}",
1081 if inflight_count == 1 { "" } else { "s" }
1082 );
1083 }
1084
1085 let mut first_error = None;
1086
1087 for (client_order_id, client_id) in cancel_routes {
1088 if let Err(e) = self.cancel_order(client_order_id, client_id, params.clone()) {
1089 if first_error.is_none() {
1090 first_error = Some(e);
1091 } else {
1092 log::error!("Error canceling {client_order_id}: {e}");
1093 }
1094 }
1095 }
1096
1097 first_error.map_or(Ok(()), Err)
1098 }
1099
1100 #[expect(clippy::too_many_arguments)]
1106 fn close_position(
1107 &mut self,
1108 position: &Position,
1109 client_id: Option<ClientId>,
1110 tags: Option<Vec<Ustr>>,
1111 time_in_force: Option<TimeInForce>,
1112 reduce_only: Option<bool>,
1113 quote_quantity: Option<bool>,
1114 params: Option<Params>,
1115 ) -> anyhow::Result<()>
1116 where
1117 Self: StrategyNative,
1118 {
1119 let core = StrategyNative::strategy_core_mut(self);
1120
1121 if position.is_closed() {
1122 log::warn!("Cannot close position (already closed): {}", position.id);
1123 return Ok(());
1124 }
1125
1126 let Some(closing_side) = OrderCore::closing_side(position.side) else {
1127 log::warn!("Cannot close flat position: {}", position.id);
1128 return Ok(());
1129 };
1130
1131 let order = core.order_factory().market(
1132 position.instrument_id,
1133 closing_side,
1134 position.quantity,
1135 time_in_force,
1136 reduce_only.or(Some(true)),
1137 quote_quantity,
1138 None,
1139 None,
1140 tags,
1141 None,
1142 );
1143
1144 self.submit_order(order, Some(position.id), client_id, params)
1145 }
1146
1147 #[expect(clippy::too_many_arguments)]
1153 fn close_all_positions(
1154 &mut self,
1155 instrument_id: InstrumentId,
1156 position_side: Option<PositionSide>,
1157 client_id: Option<ClientId>,
1158 tags: Option<Vec<Ustr>>,
1159 time_in_force: Option<TimeInForce>,
1160 reduce_only: Option<bool>,
1161 quote_quantity: Option<bool>,
1162 params: Option<Params>,
1163 ) -> anyhow::Result<()>
1164 where
1165 Self: StrategyNative,
1166 {
1167 let core = StrategyNative::strategy_core_mut(self);
1168 let strategy_id = registered_strategy_id(core)?;
1169 let cache = core.cache_ref();
1170
1171 let positions_open = cache.positions_open(
1172 None,
1173 Some(&instrument_id),
1174 Some(&strategy_id),
1175 None,
1176 position_side,
1177 );
1178
1179 let side_str = position_side.map(|s| format!(" {s}")).unwrap_or_default();
1180
1181 if positions_open.is_empty() {
1182 log::info!("No {instrument_id} open{side_str} positions to close");
1183 return Ok(());
1184 }
1185
1186 let count = positions_open.len();
1187 log::info!(
1188 "Closing {count} open{side_str} position{}",
1189 if count == 1 { "" } else { "s" }
1190 );
1191
1192 let positions_data: Vec<_> = positions_open
1193 .iter()
1194 .map(|p| (p.id, p.instrument_id, p.side, p.quantity, p.is_closed()))
1195 .collect();
1196 drop(positions_open);
1197
1198 drop(cache);
1199
1200 for (pos_id, pos_instrument_id, pos_side, pos_quantity, is_closed) in positions_data {
1201 if is_closed {
1202 continue;
1203 }
1204
1205 let core = StrategyNative::strategy_core_mut(self);
1206 let Some(closing_side) = OrderCore::closing_side(pos_side) else {
1207 continue;
1208 };
1209 let order = core.order_factory().market(
1210 pos_instrument_id,
1211 closing_side,
1212 pos_quantity,
1213 time_in_force,
1214 reduce_only.or(Some(true)),
1215 quote_quantity,
1216 None,
1217 None,
1218 tags.clone(),
1219 None,
1220 );
1221
1222 self.submit_order(order, Some(pos_id), client_id, params.clone())?;
1223 }
1224
1225 Ok(())
1226 }
1227
1228 fn query_account(
1237 &mut self,
1238 account_id: AccountId,
1239 client_id: Option<ClientId>,
1240 params: Option<Params>,
1241 ) -> anyhow::Result<()>
1242 where
1243 Self: StrategyNative,
1244 {
1245 let core = StrategyNative::strategy_core_mut(self);
1246
1247 let trader_id = registered_trader_id(core)?;
1248 let ts_init = core.clock_mut().timestamp_ns();
1249
1250 let command = QueryAccount::new(
1251 trader_id,
1252 client_id,
1253 account_id,
1254 UUID4::new(),
1255 ts_init,
1256 params,
1257 None, );
1259
1260 send_exec_command(TradingCommand::QueryAccount(command));
1261 Ok(())
1262 }
1263
1264 fn query_order(
1273 &mut self,
1274 order: &OrderAny,
1275 client_id: Option<ClientId>,
1276 params: Option<Params>,
1277 ) -> anyhow::Result<()>
1278 where
1279 Self: StrategyNative,
1280 {
1281 let core = StrategyNative::strategy_core_mut(self);
1282
1283 let trader_id = registered_trader_id(core)?;
1284 let strategy_id = registered_strategy_id(core)?;
1285 let ts_init = core.clock_mut().timestamp_ns();
1286
1287 let command = QueryOrder::new(
1288 trader_id,
1289 client_id,
1290 strategy_id,
1291 order.instrument_id(),
1292 order.client_order_id(),
1293 order.venue_order_id(),
1294 UUID4::new(),
1295 ts_init,
1296 params,
1297 None, );
1299
1300 send_exec_command(TradingCommand::QueryOrder(command));
1301 Ok(())
1302 }
1303
1304 fn handle_order_event(&mut self, event: OrderEventAny)
1306 where
1307 Self: StrategyNative,
1308 {
1309 let state = {
1310 let core = StrategyNative::strategy_core_mut(self);
1311 let id = &core.actor.actor_id;
1312 let is_warning = matches!(
1313 &event,
1314 OrderEventAny::Denied(_)
1315 | OrderEventAny::Rejected(_)
1316 | OrderEventAny::CancelRejected(_)
1317 | OrderEventAny::ModifyRejected(_)
1318 );
1319
1320 if is_warning {
1321 log::warn!("{id} {RECV}{EVT} {event}");
1322 } else if core.actor.config.log_events {
1323 log::info!("{id} {RECV}{EVT} {event}");
1324 }
1325
1326 core.actor.state()
1327 };
1328
1329 let client_order_id = event.client_order_id();
1330 let cached_order_is_closed = {
1331 let core = StrategyNative::strategy_core_mut(self);
1332 core.cache_ref()
1333 .order(&client_order_id)
1334 .map(|order| order.is_closed())
1335 };
1336 let is_terminal = match &event {
1337 OrderEventAny::FillVoided(event) => {
1338 cached_order_is_closed.unwrap_or(!event.is_reopened)
1339 }
1340 OrderEventAny::Filled(_) => cached_order_is_closed.unwrap_or(true),
1341 OrderEventAny::Canceled(_)
1342 | OrderEventAny::Rejected(_)
1343 | OrderEventAny::Expired(_)
1344 | OrderEventAny::Denied(_) => true,
1345 _ => false,
1346 };
1347
1348 if is_terminal {
1351 self.cancel_gtd_expiry(&client_order_id);
1352 }
1353
1354 if state != ComponentState::Running {
1357 return;
1358 }
1359
1360 if matches!(&event, OrderEventAny::CancelRejected(_)) {
1361 let order = StrategyNative::strategy_core_mut(self)
1362 .cache_ref()
1363 .order(&client_order_id)
1364 .map(|order| order.clone());
1365 if let Some(order) = order
1366 && (order.is_open() || order.is_inflight())
1367 && !self.has_gtd_expiry_timer(&client_order_id)
1368 && let Err(e) = self.set_gtd_expiry(&order)
1369 {
1370 log::error!(
1371 "Failed to restore GTD expiry for cancel-rejected order {client_order_id}: {e}"
1372 );
1373 }
1374 }
1375
1376 if matches!(&event, OrderEventAny::FillVoided(event) if event.is_reopened) {
1377 let order = StrategyNative::strategy_core_mut(self)
1378 .cache_ref()
1379 .order(&client_order_id)
1380 .map(|order| order.clone());
1381 if let Some(order) = order
1382 && order.is_open()
1383 && !self.has_gtd_expiry_timer(&client_order_id)
1384 && let Err(e) = self.set_gtd_expiry(&order)
1385 {
1386 log::error!(
1387 "Failed to restore GTD expiry for reopened order {client_order_id}: {e}"
1388 );
1389 }
1390 }
1391
1392 let manager_actions = {
1393 let core = StrategyNative::strategy_core_mut(self);
1394 if core.config.manage_contingent_orders {
1395 core.order_manager
1396 .as_mut()
1397 .map_or_else(Vec::new, |manager| manager.handle_event(&event))
1398 } else {
1399 Vec::new()
1400 }
1401 };
1402 self.dispatch_manager_actions(manager_actions);
1403
1404 match &event {
1405 OrderEventAny::Initialized(e) => self.on_order_initialized(e.clone()),
1406 OrderEventAny::Denied(e) => self.on_order_denied(*e),
1407 OrderEventAny::Emulated(e) => self.on_order_emulated(*e),
1408 OrderEventAny::Released(e) => self.on_order_released(*e),
1409 OrderEventAny::Submitted(e) => self.on_order_submitted(*e),
1410 OrderEventAny::Rejected(e) => self.on_order_rejected(*e),
1411 OrderEventAny::Accepted(e) => self.on_order_accepted(*e),
1412 OrderEventAny::Canceled(e) => self.on_order_canceled(e),
1413 OrderEventAny::Expired(e) => self.on_order_expired(*e),
1414 OrderEventAny::Triggered(e) => self.on_order_triggered(*e),
1415 OrderEventAny::PendingUpdate(e) => self.on_order_pending_update(*e),
1416 OrderEventAny::PendingCancel(e) => self.on_order_pending_cancel(*e),
1417 OrderEventAny::ModifyRejected(e) => self.on_order_modify_rejected(*e),
1418 OrderEventAny::CancelRejected(e) => self.on_order_cancel_rejected(*e),
1419 OrderEventAny::Updated(e) => self.on_order_updated(*e),
1420 OrderEventAny::Filled(e) => self.on_order_filled(e),
1421 OrderEventAny::FillVoided(e) => self.on_order_fill_voided(e),
1422 }
1423 self.on_order_event(event);
1424 }
1425
1426 fn dispatch_manager_actions(&mut self, actions: Vec<OrderManagerAction>)
1427 where
1428 Self: StrategyNative,
1429 {
1430 for action in actions {
1431 match action {
1432 OrderManagerAction::PublishInitialized(event) => {
1433 let topic = msgbus::switchboard::get_event_order_topic(event.strategy_id());
1434 msgbus::publish_order_event(topic, &event);
1435 }
1436 OrderManagerAction::SubmitToEmulator(command) => {
1437 send_emulator_command(TradingCommand::SubmitOrder(command));
1438 }
1439 OrderManagerAction::SubmitToRisk(command) => {
1440 send_risk_command(TradingCommand::SubmitOrder(command));
1441 }
1442 OrderManagerAction::SubmitToAlgorithm {
1443 command,
1444 exec_algorithm_id,
1445 } => send_algo_command(command, exec_algorithm_id),
1446 OrderManagerAction::CancelLocal(order) => {
1447 let client_order_id = order.client_order_id();
1448 if let Err(e) = self.cancel_order(client_order_id, None, None) {
1449 log::error!(
1450 "Failed to dispatch contingent cancel for {client_order_id}: {e}"
1451 );
1452 }
1453 }
1454 OrderManagerAction::ModifyLocalQuantity { order, quantity } => {
1455 let client_order_id = order.client_order_id();
1456 if let Err(e) =
1457 self.modify_order(client_order_id, Some(quantity), None, None, None, None)
1458 {
1459 log::error!(
1460 "Failed to dispatch contingent modify for {client_order_id}: {e}"
1461 );
1462 }
1463 }
1464 }
1465 }
1466 }
1467
1468 fn handle_position_event(&mut self, event: PositionEvent)
1470 where
1471 Self: StrategyNative,
1472 {
1473 let state = {
1474 let core = StrategyNative::strategy_core_mut(self);
1475
1476 if core.actor.config.log_events {
1477 let id = &core.actor.actor_id;
1478 log::info!("{id} {RECV}{EVT} {event:?}");
1479 }
1480
1481 core.actor.state()
1482 };
1483
1484 if state != ComponentState::Running {
1485 return;
1486 }
1487
1488 match &event {
1489 PositionEvent::PositionOpened(e) => self.on_position_opened(e.clone()),
1490 PositionEvent::PositionChanged(e) => self.on_position_changed(e.clone()),
1491 PositionEvent::PositionClosed(e) => self.on_position_closed(e.clone()),
1492 PositionEvent::PositionAdjusted(_) => {
1493 return;
1494 }
1495 }
1496 self.on_position_event(event);
1497 }
1498
1499 fn on_start(&mut self) -> anyhow::Result<()>
1510 where
1511 Self: StrategyNative,
1512 {
1513 let core = StrategyNative::strategy_core_mut(self);
1514 let strategy_id = registered_strategy_id(core)?;
1515 log::info!("Starting {strategy_id}");
1516
1517 if core.config.manage_gtd_expiry {
1518 self.reactivate_gtd_timers();
1519 }
1520
1521 Ok(())
1522 }
1523
1524 fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()>
1534 where
1535 Self: StrategyNative + Component,
1536 {
1537 route_time_event(self, event);
1538 Ok(())
1539 }
1540
1541 #[allow(unused_variables)]
1547 fn on_order_initialized(&mut self, event: OrderInitialized) {}
1548
1549 #[allow(unused_variables)]
1553 fn on_order_event(&mut self, event: OrderEventAny) {}
1554
1555 #[allow(unused_variables)]
1559 fn on_order_denied(&mut self, event: OrderDenied) {}
1560
1561 #[allow(unused_variables)]
1565 fn on_order_emulated(&mut self, event: OrderEmulated) {}
1566
1567 #[allow(unused_variables)]
1571 fn on_order_released(&mut self, event: OrderReleased) {}
1572
1573 #[allow(unused_variables)]
1577 fn on_order_submitted(&mut self, event: OrderSubmitted) {}
1578
1579 #[allow(unused_variables)]
1583 fn on_order_rejected(&mut self, event: OrderRejected) {}
1584
1585 #[allow(unused_variables)]
1589 fn on_order_accepted(&mut self, event: OrderAccepted) {}
1590
1591 #[allow(unused_variables)]
1595 fn on_order_expired(&mut self, event: OrderExpired) {}
1596
1597 #[allow(unused_variables)]
1601 fn on_order_triggered(&mut self, event: OrderTriggered) {}
1602
1603 #[allow(unused_variables)]
1607 fn on_order_pending_update(&mut self, event: OrderPendingUpdate) {}
1608
1609 #[allow(unused_variables)]
1613 fn on_order_pending_cancel(&mut self, event: OrderPendingCancel) {}
1614
1615 #[allow(unused_variables)]
1619 fn on_order_modify_rejected(&mut self, event: OrderModifyRejected) {}
1620
1621 #[allow(unused_variables)]
1625 fn on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {}
1626
1627 #[allow(unused_variables)]
1631 fn on_order_updated(&mut self, event: OrderUpdated) {}
1632
1633 #[allow(unused_variables)]
1637 fn on_order_canceled(&mut self, event: &OrderCanceled) {}
1638
1639 #[allow(unused_variables)]
1643 fn on_order_filled(&mut self, event: &OrderFilled) {}
1644
1645 #[allow(unused_variables)]
1647 fn on_order_fill_voided(&mut self, event: &OrderFillVoided) {}
1648
1649 #[allow(unused_variables)]
1653 fn on_position_opened(&mut self, event: PositionOpened) {}
1654
1655 #[allow(unused_variables)]
1659 fn on_position_event(&mut self, event: PositionEvent) {}
1660
1661 #[allow(unused_variables)]
1665 fn on_position_changed(&mut self, event: PositionChanged) {}
1666
1667 #[allow(unused_variables)]
1671 fn on_position_closed(&mut self, event: PositionClosed) {}
1672
1673 fn on_market_exit(&mut self) {}
1677
1678 fn post_market_exit(&mut self) {}
1682
1683 fn is_exiting(&self) -> bool
1687 where
1688 Self: StrategyNative,
1689 {
1690 StrategyNative::strategy_core(self).is_exiting
1691 }
1692
1693 fn market_exit(&mut self) -> anyhow::Result<()>
1709 where
1710 Self: StrategyNative,
1711 {
1712 let core = StrategyNative::strategy_core_mut(self);
1713 let strategy_id = registered_strategy_id(core)?;
1714
1715 if core.actor.state() != ComponentState::Running {
1716 log::warn!("{strategy_id} Cannot market exit: strategy is not running");
1717 return Ok(());
1718 }
1719
1720 if core.is_exiting {
1721 log::warn!("{strategy_id} Market exit called when already in progress");
1722 return Ok(());
1723 }
1724
1725 core.is_exiting = true;
1726 core.market_exit_attempts = 0;
1727 let time_in_force = core.config.market_exit_time_in_force;
1728 let reduce_only = core.config.market_exit_reduce_only;
1729
1730 log::info!("{strategy_id} Initiating market exit...");
1731
1732 self.on_market_exit();
1733
1734 let core = StrategyNative::strategy_core_mut(self);
1735 let cache = core.cache_ref();
1736
1737 let mut instruments: AHashSet<InstrumentId> = AHashSet::new();
1738
1739 for client_order_id in
1740 cache.iter_client_order_ids_open(None, None, Some(&strategy_id), None)
1741 {
1742 if let Some(order) = cache.order(&client_order_id) {
1743 instruments.insert(order.instrument_id());
1744 }
1745 }
1746
1747 for client_order_id in
1748 cache.iter_client_order_ids_inflight(None, None, Some(&strategy_id), None)
1749 {
1750 if let Some(order) = cache.order(&client_order_id) {
1751 instruments.insert(order.instrument_id());
1752 }
1753 }
1754
1755 for position_id in cache.iter_position_open_ids(None, None, Some(&strategy_id), None) {
1756 if let Some(position) = cache.position(&position_id) {
1757 instruments.insert(position.instrument_id);
1758 }
1759 }
1760
1761 let market_exit_tag = core.market_exit_tag;
1762 let mut instruments: Vec<_> = instruments.into_iter().collect();
1766 instruments.sort();
1767 drop(cache);
1768
1769 for instrument_id in instruments {
1770 if let Err(e) = self.cancel_all_orders(instrument_id, None, None, true, None) {
1771 log::error!("Error canceling orders for {instrument_id}: {e}");
1772 }
1773
1774 if let Err(e) = self.close_all_positions(
1775 instrument_id,
1776 None,
1777 None,
1778 Some(vec![market_exit_tag]),
1779 Some(time_in_force),
1780 Some(reduce_only),
1781 None,
1782 None,
1783 ) {
1784 log::error!("Error closing positions for {instrument_id}: {e}");
1785 }
1786 }
1787
1788 let core = StrategyNative::strategy_core_mut(self);
1789 let interval_ms = core.config.market_exit_interval_ms;
1790 let timer_name = core.market_exit_timer_name;
1791
1792 log::info!("{strategy_id} Setting market exit timer at {interval_ms}ms intervals");
1793
1794 let Ok(interval_ns) = DurationNanos::try_from_millis(interval_ms) else {
1795 core.is_exiting = false;
1796 core.market_exit_attempts = 0;
1797 anyhow::bail!("Market exit timer interval exceeds the nanosecond range");
1798 };
1799 let result = core.clock_mut().set_timer_ns(
1800 timer_name.as_str(),
1801 interval_ns,
1802 None,
1803 None,
1804 None,
1805 None,
1806 None,
1807 );
1808
1809 if let Err(e) = result {
1810 core.is_exiting = false;
1812 core.market_exit_attempts = 0;
1813 return Err(e);
1814 }
1815
1816 Ok(())
1817 }
1818
1819 fn check_market_exit(&mut self, _event: TimeEvent)
1823 where
1824 Self: StrategyNative + Component,
1825 {
1826 if !self.is_exiting() {
1828 return;
1829 }
1830
1831 let core = StrategyNative::strategy_core_mut(self);
1832 let Some(strategy_id) = core.strategy_id() else {
1833 log::error!("Cannot check market exit: strategy_id is not set");
1834 return;
1835 };
1836
1837 core.market_exit_attempts += 1;
1838 let attempts = core.market_exit_attempts;
1839 let max_attempts = core.config.market_exit_max_attempts;
1840
1841 log::debug!(
1842 "{strategy_id} Market exit check triggered (attempt {attempts}/{max_attempts})"
1843 );
1844
1845 if attempts >= max_attempts {
1846 let cache = core.cache_ref();
1847 let open_orders_count =
1848 cache.orders_open_count(None, None, Some(&strategy_id), None, None);
1849 let inflight_orders_count =
1850 cache.orders_inflight_count(None, None, Some(&strategy_id), None, None);
1851 let open_positions_count =
1852 cache.positions_open_count(None, None, Some(&strategy_id), None, None);
1853
1854 drop(cache);
1855
1856 log::warn!(
1857 "{strategy_id} Market exit max attempts ({max_attempts}) reached, \
1858 completing with open orders: {open_orders_count}, \
1859 inflight orders: {inflight_orders_count}, \
1860 open positions: {open_positions_count}"
1861 );
1862
1863 self.finalize_market_exit();
1864 return;
1865 }
1866
1867 let cache = core.cache_ref();
1868 let has_open_orders = !cache
1869 .orders_open(None, None, Some(&strategy_id), None, None)
1870 .is_empty();
1871 let has_inflight_orders = !cache
1872 .orders_inflight(None, None, Some(&strategy_id), None, None)
1873 .is_empty();
1874
1875 if has_open_orders || has_inflight_orders {
1876 return;
1877 }
1878
1879 let positions_data: Vec<_> = cache
1880 .positions_open(None, None, Some(&strategy_id), None, None)
1881 .iter()
1882 .map(|p| (p.id, p.instrument_id, p.side, p.quantity, p.is_closed()))
1883 .collect();
1884
1885 if !positions_data.is_empty() {
1886 drop(cache);
1888
1889 for (pos_id, instrument_id, side, quantity, is_closed) in positions_data {
1890 if is_closed {
1891 continue;
1892 }
1893
1894 let core = StrategyNative::strategy_core_mut(self);
1895 let time_in_force = core.config.market_exit_time_in_force;
1896 let reduce_only = core.config.market_exit_reduce_only;
1897 let market_exit_tag = core.market_exit_tag;
1898 let Some(closing_side) = OrderCore::closing_side(side) else {
1899 continue;
1900 };
1901 let order = core.order_factory().market(
1902 instrument_id,
1903 closing_side,
1904 quantity,
1905 Some(time_in_force),
1906 Some(reduce_only),
1907 None,
1908 None,
1909 None,
1910 Some(vec![market_exit_tag]),
1911 None,
1912 );
1913
1914 if let Err(e) = self.submit_order(order, Some(pos_id), None, None) {
1915 log::error!("Error re-submitting close order for position {pos_id}: {e}");
1916 }
1917 }
1918 return;
1919 }
1920
1921 drop(cache);
1922 self.finalize_market_exit();
1923 }
1924
1925 fn finalize_market_exit(&mut self)
1930 where
1931 Self: StrategyNative + Component,
1932 {
1933 let (actor_id, should_stop) = {
1934 let core = StrategyNative::strategy_core_mut(self);
1935 let actor_id = core.actor_id();
1936 let should_stop = core.pending_stop;
1937 (actor_id, should_stop)
1938 };
1939
1940 self.cancel_market_exit();
1941
1942 let hook_result = catch_unwind(AssertUnwindSafe(|| {
1943 self.post_market_exit();
1944 }));
1945
1946 if let Err(e) = hook_result {
1947 log::error!("{actor_id} Error in post_market_exit: {e:?}");
1948 }
1949
1950 if should_stop {
1951 log::info!("{actor_id} Market exit complete, stopping strategy");
1952
1953 if let Err(e) = Component::stop(self) {
1954 log::error!("{actor_id} Failed to stop: {e}");
1955 }
1956 }
1957
1958 let core = StrategyNative::strategy_core_mut(self);
1959 debug_assert!(
1960 !(core.pending_stop
1961 && !core.is_exiting
1962 && core.actor.state() == ComponentState::Running),
1963 "INVARIANT: stuck state after finalize_market_exit"
1964 );
1965 }
1966
1967 fn cancel_market_exit(&mut self)
1971 where
1972 Self: StrategyNative,
1973 {
1974 let core = StrategyNative::strategy_core_mut(self);
1975 let timer_name = core.market_exit_timer_name;
1976
1977 if core
1978 .clock_mut()
1979 .timer_names()
1980 .contains(&timer_name.as_str())
1981 {
1982 core.clock_mut().cancel_timer(timer_name.as_str());
1983 }
1984
1985 core.is_exiting = false;
1986 core.pending_stop = false;
1987 core.market_exit_attempts = 0;
1988 }
1989
1990 fn stop(&mut self) -> bool
2002 where
2003 Self: StrategyNative,
2004 {
2005 let (manage_stop, is_exiting, should_initiate_exit) = {
2006 let core = StrategyNative::strategy_core_mut(self);
2007 let actor_id = core.actor_id();
2008 let manage_stop = core.config.manage_stop;
2009 let state = core.actor.state();
2010 let pending_stop = core.pending_stop;
2011 let is_exiting = core.is_exiting;
2012
2013 if manage_stop {
2014 if state != ComponentState::Running {
2015 return true; }
2017
2018 if pending_stop {
2019 return false; }
2021
2022 core.pending_stop = true;
2023 let should_initiate_exit = !is_exiting;
2024
2025 if should_initiate_exit {
2026 log::info!("{actor_id} Initiating market exit before stop");
2027 }
2028
2029 (manage_stop, is_exiting, should_initiate_exit)
2030 } else {
2031 (manage_stop, is_exiting, false)
2032 }
2033 };
2034
2035 if manage_stop {
2036 if should_initiate_exit && let Err(e) = self.market_exit() {
2037 log::warn!("Market exit failed during stop: {e}, proceeding with stop");
2038 StrategyNative::strategy_core_mut(self).pending_stop = false;
2039 return true;
2040 }
2041 debug_assert!(
2042 self.is_exiting(),
2043 "INVARIANT: deferring stop but not exiting"
2044 );
2045 return false; }
2047
2048 if is_exiting {
2050 self.cancel_market_exit();
2051 }
2052
2053 true }
2055
2056 fn deny_order(&mut self, order: &OrderAny, reason: Ustr)
2061 where
2062 Self: StrategyNative,
2063 {
2064 let core = StrategyNative::strategy_core_mut(self);
2065 let Some(trader_id) = core.trader_id() else {
2066 log::error!(
2067 "Cannot deny order {}: trader_id is not set",
2068 order.client_order_id()
2069 );
2070 return;
2071 };
2072 let Some(strategy_id) = core.strategy_id() else {
2073 log::error!(
2074 "Cannot deny order {}: strategy_id is not set",
2075 order.client_order_id()
2076 );
2077 return;
2078 };
2079 let ts_now = core.clock_mut().timestamp_ns();
2080
2081 let event = OrderDenied::new(
2082 trader_id,
2083 strategy_id,
2084 order.instrument_id(),
2085 order.client_order_id(),
2086 reason,
2087 UUID4::new(),
2088 ts_now,
2089 ts_now,
2090 );
2091
2092 log::warn!(
2093 "{strategy_id} Order {} denied: {reason}",
2094 order.client_order_id()
2095 );
2096
2097 let publish_initialized = {
2098 let cache_rc = core.cache_rc();
2099 let mut cache = cache_rc.borrow_mut();
2100 if cache.order_exists(&order.client_order_id()) {
2101 false
2102 } else {
2103 match cache.add_order(order.clone(), None, None, true) {
2104 Ok(()) => true,
2105 Err(e) => {
2106 log::warn!("Failed to add denied order to cache: {e}");
2107 false
2108 }
2109 }
2110 }
2111 };
2112
2113 if publish_initialized {
2114 publish_order_initialized(order);
2115 }
2116
2117 let event = OrderEventAny::Denied(event);
2118 let applied = {
2119 let cache_rc = core.cache_rc();
2120 let mut cache = cache_rc.borrow_mut();
2121 if let Err(e) = cache.update_order(&event) {
2122 log::warn!("Failed to apply OrderDenied event: {e}");
2123 false
2124 } else {
2125 true
2126 }
2127 };
2128
2129 if applied {
2130 let topic = format!("events.order.{strategy_id}");
2131 msgbus::publish_order_event(topic.into(), &event);
2132 }
2133 }
2134
2135 fn deny_order_list(&mut self, orders: &[OrderAny], reason: Ustr)
2139 where
2140 Self: StrategyNative,
2141 {
2142 for order in orders {
2143 if !order.is_closed() {
2144 self.deny_order(order, reason);
2145 }
2146 }
2147 }
2148
2149 fn set_gtd_expiry(&mut self, order: &OrderAny) -> anyhow::Result<()>
2159 where
2160 Self: StrategyNative,
2161 {
2162 let core = StrategyNative::strategy_core_mut(self);
2163
2164 if !core.config.manage_gtd_expiry || order.time_in_force() != TimeInForce::Gtd {
2165 return Ok(());
2166 }
2167
2168 let Some(expire_time) = order.expire_time() else {
2169 return Ok(());
2170 };
2171
2172 let client_order_id = order.client_order_id();
2173 let timer_name = format!("GTD-EXPIRY:{client_order_id}");
2174
2175 let current_time_ns = {
2176 let clock = core.clock_mut();
2177 clock.timestamp_ns()
2178 };
2179
2180 if current_time_ns >= expire_time.as_u64() {
2181 log::info!("GTD order {client_order_id} already expired, canceling immediately");
2182 return self.cancel_order(order.client_order_id(), None, None);
2183 }
2184
2185 {
2186 let mut clock = core.clock_mut();
2187 clock.set_time_alert_ns(&timer_name, expire_time, None, None)?;
2188 }
2189
2190 core.gtd_timers
2191 .insert(client_order_id, Ustr::from(&timer_name));
2192
2193 log::debug!("Set GTD expiry timer for {client_order_id} at {expire_time}");
2194 Ok(())
2195 }
2196
2197 fn cancel_gtd_expiry(&mut self, client_order_id: &ClientOrderId)
2199 where
2200 Self: StrategyNative,
2201 {
2202 let core = StrategyNative::strategy_core_mut(self);
2203
2204 if let Some(timer_name) = core.gtd_timers.remove(client_order_id) {
2205 core.clock_mut().cancel_timer(timer_name.as_str());
2206 log::debug!("Canceled GTD expiry timer for {client_order_id}");
2207 }
2208 }
2209
2210 fn has_gtd_expiry_timer(&mut self, client_order_id: &ClientOrderId) -> bool
2212 where
2213 Self: StrategyNative,
2214 {
2215 let core = StrategyNative::strategy_core_mut(self);
2216 core.gtd_timers.contains_key(client_order_id)
2217 }
2218
2219 fn expire_gtd_order(&mut self, event: TimeEvent)
2223 where
2224 Self: StrategyNative,
2225 {
2226 let timer_name = event.name;
2227 let Some(client_order_id) = timer_name
2228 .strip_prefix("GTD-EXPIRY:")
2229 .and_then(|value| ClientOrderId::new_checked(value).ok())
2230 else {
2231 log::error!("Invalid GTD timer name format: {timer_name}");
2232 return;
2233 };
2234
2235 let core = StrategyNative::strategy_core_mut(self);
2236 if core.gtd_timers.get(&client_order_id) != Some(&timer_name) {
2237 return;
2238 }
2239 core.gtd_timers.remove(&client_order_id);
2240
2241 let order = core.cache_ref().order(&client_order_id).map(|o| o.clone());
2242 let Some(order) = order else {
2243 log::warn!("GTD order {client_order_id} not found in cache");
2244 return;
2245 };
2246
2247 log::info!("GTD order {client_order_id} expired");
2248
2249 if let Err(e) = self.cancel_order(order.client_order_id(), None, None) {
2250 log::error!("Failed to cancel expired GTD order {client_order_id}: {e}");
2251 }
2252 }
2253
2254 fn reactivate_gtd_timers(&mut self)
2259 where
2260 Self: StrategyNative,
2261 {
2262 let core = StrategyNative::strategy_core_mut(self);
2263 let Some(strategy_id) = core.strategy_id() else {
2264 log::error!("Cannot reactivate GTD timers: strategy_id is not set");
2265 return;
2266 };
2267 let current_time_ns = core.clock_mut().timestamp_ns();
2268
2269 let gtd_orders: Vec<OrderAny> = core
2270 .cache_ref()
2271 .orders_open(None, None, Some(&strategy_id), None, None)
2272 .into_iter()
2273 .filter(|o| o.time_in_force() == TimeInForce::Gtd)
2274 .map(|o| o.clone())
2275 .collect();
2276
2277 for order in gtd_orders {
2278 let Some(expire_time) = order.expire_time() else {
2279 continue;
2280 };
2281
2282 let expire_time_ns = expire_time.as_u64();
2283 let client_order_id = order.client_order_id();
2284
2285 if current_time_ns >= expire_time_ns {
2286 log::info!("GTD order {client_order_id} already expired, canceling immediately");
2287 if let Err(e) = self.cancel_order(order.client_order_id(), None, None) {
2288 log::error!("Failed to cancel expired GTD order {client_order_id}: {e}");
2289 }
2290 } else if let Err(e) = self.set_gtd_expiry(&order) {
2291 log::error!("Failed to set GTD expiry timer for {client_order_id}: {e}");
2292 }
2293 }
2294 }
2295}
2296
2297pub fn route_time_event<T>(strategy: &mut T, event: &TimeEvent)
2303where
2304 T: Strategy + StrategyNative + Component + ?Sized,
2305{
2306 let (gtd_order_id, is_market_exit) = {
2307 let core = StrategyNative::strategy_core(strategy);
2308 let gtd_order_id = event
2309 .name
2310 .strip_prefix("GTD-EXPIRY:")
2311 .and_then(|value| ClientOrderId::new_checked(value).ok())
2312 .filter(|client_order_id| core.gtd_timers.get(client_order_id) == Some(&event.name));
2313 let is_market_exit = event.name == core.market_exit_timer_name;
2314 (gtd_order_id, is_market_exit)
2315 };
2316
2317 if gtd_order_id.is_none() && !is_market_exit {
2318 return;
2319 }
2320
2321 let core = StrategyNative::strategy_core_mut(strategy);
2322 if core.managed_time_event_last_id == Some(event.event_id) {
2323 return;
2324 }
2325 core.managed_time_event_last_id = Some(event.event_id);
2326
2327 if gtd_order_id.is_some() {
2328 strategy.expire_gtd_order(event.clone());
2329 } else {
2330 strategy.check_market_exit(event.clone());
2331 }
2332}
2333
2334pub(super) fn submit_order_native<T>(
2335 strategy: &mut T,
2336 order: &OrderAny,
2337 position_id: Option<PositionId>,
2338 client_id: Option<ClientId>,
2339 params: Option<Params>,
2340) -> anyhow::Result<()>
2341where
2342 T: Strategy + StrategyNative + ?Sized,
2343{
2344 let core = StrategyNative::strategy_core_mut(strategy);
2345
2346 let trader_id = registered_trader_id(core)?;
2347 let strategy_id = registered_strategy_id(core)?;
2348 let ts_init = core.clock_mut().timestamp_ns();
2349
2350 if order.status() != OrderStatus::Initialized {
2351 anyhow::bail!(
2352 "Order denied: invalid status for {}, expected INITIALIZED",
2353 order.client_order_id()
2354 );
2355 }
2356
2357 let market_exit_tag = core.market_exit_tag;
2358 let is_market_exit_order = order
2359 .tags()
2360 .is_some_and(|tags| tags.contains(&market_exit_tag));
2361 let should_deny_for_market_exit =
2362 core.is_exiting && !order.is_reduce_only() && !is_market_exit_order;
2363
2364 if should_deny_for_market_exit {
2365 strategy.deny_order(order, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
2366 return Ok(());
2367 }
2368
2369 let core = StrategyNative::strategy_core_mut(strategy);
2370 let params = params.filter(|params| !params.is_empty());
2371
2372 {
2373 let cache_rc = core.cache_rc();
2374 let mut cache = cache_rc.try_borrow_mut().map_err(|_| {
2375 anyhow::anyhow!(
2376 "Cannot submit order {}: cache is currently borrowed",
2377 order.client_order_id()
2378 )
2379 })?;
2380 cache.add_order(order.clone(), position_id, client_id, true)?;
2381 }
2382
2383 publish_order_initialized(order);
2384
2385 let command = SubmitOrder::new(
2386 trader_id,
2387 client_id,
2388 strategy_id,
2389 order.instrument_id(),
2390 order.client_order_id(),
2391 order.init_event().clone(),
2392 order.exec_algorithm_id(),
2393 position_id,
2394 params,
2395 UUID4::new(),
2396 ts_init,
2397 None, );
2399
2400 if order.emulation_trigger().is_some() {
2401 send_emulator_command(TradingCommand::SubmitOrder(command));
2402 } else if let Some(exec_algorithm_id) = order.exec_algorithm_id() {
2403 send_algo_command(command, exec_algorithm_id);
2404 } else {
2405 send_risk_command(TradingCommand::SubmitOrder(command));
2406 }
2407
2408 strategy.set_gtd_expiry(order)?;
2409 Ok(())
2410}
2411
2412fn publish_order_initialized(order: &OrderAny) {
2413 let topic = format!("events.order.{}", order.strategy_id());
2414 let event = OrderEventAny::Initialized(order.init_event().clone());
2415 msgbus::publish_order_event(topic.into(), &event);
2416}
2417
2418fn send_emulator_command(command: TradingCommand) {
2419 log_cmd_send(&command);
2420 let endpoint = MessagingSwitchboard::order_emulator_execute();
2421 msgbus::send_trading_command(endpoint, command);
2422}
2423
2424fn send_algo_command(command: SubmitOrder, exec_algorithm_id: ExecAlgorithmId) {
2425 let id = command.strategy_id;
2426 log::info!("{id} {CMD}{SEND} {command}");
2427
2428 let endpoint = format!("{exec_algorithm_id}.execute");
2429 msgbus::send_any(endpoint.into(), &TradingCommand::SubmitOrder(command));
2430}
2431
2432fn send_risk_command(command: TradingCommand) {
2433 log_cmd_send(&command);
2434 let endpoint = MessagingSwitchboard::risk_engine_queue_execute();
2435 msgbus::send_trading_command(endpoint, command);
2436}
2437
2438fn send_exec_command(command: TradingCommand) {
2439 log_cmd_send(&command);
2440 let endpoint = MessagingSwitchboard::exec_engine_queue_execute();
2441 msgbus::send_trading_command(endpoint, command);
2442}
2443
2444fn log_cmd_send(command: &TradingCommand) {
2445 if let Some(id) = command.strategy_id() {
2446 log::info!("{id} {CMD}{SEND} {command}");
2447 } else {
2448 log::info!("{CMD}{SEND} {command}");
2449 }
2450}
2451
2452fn registered_trader_id(core: &StrategyCore) -> anyhow::Result<TraderId> {
2453 core.trader_id()
2454 .ok_or_else(|| anyhow::anyhow!("Strategy not registered: trader_id is not set"))
2455}
2456
2457fn registered_strategy_id(core: &StrategyCore) -> anyhow::Result<StrategyId> {
2458 core.strategy_id()
2459 .ok_or_else(|| anyhow::anyhow!("Strategy not registered: strategy_id is not set"))
2460}
2461
2462fn required_account_id(order: &OrderAny, operation: &str) -> anyhow::Result<AccountId> {
2463 order.account_id().ok_or_else(|| {
2464 anyhow::anyhow!(
2465 "Cannot generate {operation} event for {}: account_id is not set",
2466 order.client_order_id()
2467 )
2468 })
2469}
2470
2471#[cfg(test)]
2472mod tests {
2473 use std::{cell::RefCell, rc::Rc};
2474
2475 use nautilus_common::{
2476 actor::{
2477 DataActor,
2478 registry::{deregister_actor, try_get_actor_unchecked},
2479 },
2480 cache::{Cache, ORDER_NOT_FOUND},
2481 clock::{Clock, TestClock},
2482 component::{Component, deregister_component, register_component_actor},
2483 enums::ComponentState,
2484 msgbus::{
2485 self, MessagingSwitchboard, TypedHandler, TypedIntoHandler,
2486 stubs::{
2487 TypedIntoMessageSavingHandler, TypedMessageSavingHandler, get_any_saving_handler,
2488 get_typed_into_message_saving_handler, get_typed_message_saving_handler,
2489 },
2490 },
2491 timer::{TimeEvent, TimeEventCallback},
2492 };
2493 use nautilus_core::{DurationNanos, UnixNanos};
2494 use nautilus_model::{
2495 enums::{
2496 ContingencyType, LiquiditySide, OrderSide, OrderStatus, OrderType,
2497 PositionAdjustmentType, PositionSide, TriggerType,
2498 },
2499 events::{
2500 OrderAccepted, OrderCanceled, OrderFilled, OrderRejected, PositionAdjusted,
2501 order::spec::{
2502 OrderAcceptedSpec, OrderCancelRejectedSpec, OrderCanceledSpec, OrderEmulatedSpec,
2503 OrderExpiredSpec, OrderFillVoidedSpec, OrderFilledSpec, OrderRejectedSpec,
2504 },
2505 },
2506 identifiers::{
2507 AccountId, ActorId, ClientOrderId, InstrumentId, OrderListId, PositionId, StrategyId,
2508 TradeId, TraderId, VenueOrderId,
2509 },
2510 orderbook::own::OwnOrderBook,
2511 orders::{LimitOrder, MarketOrder, OrderTestBuilder, stubs::TestOrderEventStubs},
2512 stubs::TestDefault,
2513 types::{Currency, Money, Price},
2514 };
2515 use nautilus_portfolio::portfolio::Portfolio;
2516 use rstest::rstest;
2517 use serde_json::Value;
2518
2519 use super::*;
2520 use crate::nautilus_strategy;
2521
2522 #[derive(Debug)]
2523 struct TestStrategy {
2524 core: StrategyCore,
2525 on_order_rejected_called: bool,
2526 on_order_event_called: bool,
2527 on_order_accepted_called: bool,
2528 on_order_canceled_called: bool,
2529 on_order_filled_called: bool,
2530 on_order_fill_voided_called: bool,
2531 on_order_expired_called: bool,
2532 order_event_timeline: Rc<RefCell<Vec<&'static str>>>,
2533 on_position_event_called: bool,
2534 on_position_opened_called: bool,
2535 on_position_changed_called: bool,
2536 on_position_closed_called: bool,
2537 }
2538
2539 #[derive(Debug)]
2540 struct CoreFreeStrategy {
2541 started: bool,
2542 }
2543
2544 #[derive(Debug)]
2545 struct InitializedModifyStrategy {
2546 core: StrategyCore,
2547 modified_quantity: Quantity,
2548 }
2549
2550 #[derive(Debug)]
2551 struct TimerOverrideStrategy {
2552 core: StrategyCore,
2553 gtd_expiries: usize,
2554 market_exit_checks: usize,
2555 }
2556
2557 impl DataActor for CoreFreeStrategy {
2558 fn on_start(&mut self) -> anyhow::Result<()> {
2559 self.started = true;
2560 Ok(())
2561 }
2562 }
2563
2564 impl Strategy for CoreFreeStrategy {}
2565
2566 impl DataActor for InitializedModifyStrategy {}
2567
2568 nautilus_strategy!(InitializedModifyStrategy, {
2569 fn on_order_initialized(&mut self, event: OrderInitialized) {
2570 self.modify_order(
2571 event.client_order_id,
2572 Some(self.modified_quantity),
2573 None,
2574 None,
2575 None,
2576 None,
2577 )
2578 .unwrap();
2579 }
2580 });
2581
2582 impl DataActor for TimerOverrideStrategy {
2583 fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()> {
2584 Strategy::on_time_event(self, event)
2585 }
2586 }
2587
2588 nautilus_strategy!(TimerOverrideStrategy, {
2589 fn check_market_exit(&mut self, _event: TimeEvent) {
2590 self.market_exit_checks += 1;
2591 }
2592
2593 fn expire_gtd_order(&mut self, _event: TimeEvent) {
2594 self.gtd_expiries += 1;
2595 }
2596 });
2597
2598 impl TestStrategy {
2599 fn new(config: StrategyConfig) -> Self {
2600 Self {
2601 core: StrategyCore::new(config),
2602 on_order_rejected_called: false,
2603 on_order_event_called: false,
2604 on_order_accepted_called: false,
2605 on_order_canceled_called: false,
2606 on_order_filled_called: false,
2607 on_order_fill_voided_called: false,
2608 on_order_expired_called: false,
2609 order_event_timeline: Rc::new(RefCell::new(Vec::new())),
2610 on_position_event_called: false,
2611 on_position_opened_called: false,
2612 on_position_changed_called: false,
2613 on_position_closed_called: false,
2614 }
2615 }
2616 }
2617
2618 impl DataActor for TestStrategy {}
2619
2620 nautilus_strategy!(TestStrategy, {
2621 fn on_order_canceled(&mut self, _event: &OrderCanceled) {
2622 self.on_order_canceled_called = true;
2623 }
2624
2625 fn on_order_filled(&mut self, _event: &OrderFilled) {
2626 self.on_order_filled_called = true;
2627 self.order_event_timeline.borrow_mut().push("specific");
2628 }
2629
2630 fn on_order_fill_voided(&mut self, _event: &OrderFillVoided) {
2631 self.on_order_fill_voided_called = true;
2632 }
2633
2634 fn on_order_rejected(&mut self, _event: OrderRejected) {
2635 self.on_order_rejected_called = true;
2636 }
2637
2638 fn on_order_event(&mut self, _event: OrderEventAny) {
2639 self.on_order_event_called = true;
2640 self.order_event_timeline.borrow_mut().push("aggregate");
2641 }
2642
2643 fn on_order_accepted(&mut self, _event: OrderAccepted) {
2644 self.on_order_accepted_called = true;
2645 }
2646
2647 fn on_order_expired(&mut self, _event: OrderExpired) {
2648 self.on_order_expired_called = true;
2649 }
2650
2651 fn on_position_opened(&mut self, _event: PositionOpened) {
2652 self.on_position_opened_called = true;
2653 }
2654
2655 fn on_position_event(&mut self, _event: PositionEvent) {
2656 self.on_position_event_called = true;
2657 }
2658
2659 fn on_position_changed(&mut self, _event: PositionChanged) {
2660 self.on_position_changed_called = true;
2661 }
2662
2663 fn on_position_closed(&mut self, _event: PositionClosed) {
2664 self.on_position_closed_called = true;
2665 }
2666 });
2667
2668 fn create_test_strategy() -> TestStrategy {
2669 let config = StrategyConfig {
2670 strategy_id: Some(StrategyId::from("TEST-001")),
2671 order_id_tag: Some("001".to_string()),
2672 ..Default::default()
2673 };
2674 TestStrategy::new(config)
2675 }
2676
2677 fn create_gtd_managed_strategy() -> TestStrategy {
2678 TestStrategy::new(StrategyConfig {
2679 strategy_id: Some(StrategyId::from("TEST-001")),
2680 order_id_tag: Some("001".to_string()),
2681 manage_gtd_expiry: true,
2682 ..Default::default()
2683 })
2684 }
2685
2686 fn register_strategy(strategy: &mut TestStrategy) {
2687 let _clock = register_strategy_with_clock(strategy);
2688 }
2689
2690 fn register_strategy_with_clock(strategy: &mut TestStrategy) -> Rc<RefCell<TestClock>> {
2691 let trader_id = TraderId::from("TRADER-001");
2692 let clock = Rc::new(RefCell::new(TestClock::new()));
2693 let cache = Rc::new(RefCell::new(Cache::default()));
2694 let portfolio = Rc::new(RefCell::new(Portfolio::new(
2695 clock.clone(),
2696 cache.clone(),
2697 None,
2698 )));
2699
2700 strategy
2701 .core
2702 .register(trader_id, clock.clone(), cache, portfolio)
2703 .unwrap();
2704 strategy.initialize().unwrap();
2705 clock
2706 }
2707
2708 fn register_gtd_strategy(strategy: &mut TestStrategy) -> Rc<RefCell<TestClock>> {
2709 let clock = register_strategy_with_clock(strategy);
2710 clock
2711 .borrow_mut()
2712 .register_default_handler(TimeEventCallback::from(|_event: TimeEvent| {}));
2713 clock
2714 }
2715
2716 fn start_strategy(strategy: &mut TestStrategy) {
2717 strategy.start().unwrap();
2718 }
2719
2720 fn stop_strategy(strategy: &mut TestStrategy) {
2721 Component::stop(strategy).unwrap();
2722 }
2723
2724 fn make_filled(client_order_id: ClientOrderId) -> OrderEventAny {
2725 OrderEventAny::Filled(
2726 OrderFilledSpec::builder()
2727 .trader_id(TraderId::from("TRADER-001"))
2728 .strategy_id(StrategyId::from("TEST-001"))
2729 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2730 .client_order_id(client_order_id)
2731 .venue_order_id(VenueOrderId::test_default())
2732 .account_id(AccountId::from("ACC-001"))
2733 .trade_id(TradeId::test_default())
2734 .last_qty(Quantity::default())
2735 .last_px(Price::default())
2736 .currency(Currency::from("USD"))
2737 .liquidity_side(LiquiditySide::Taker)
2738 .event_id(UUID4::default())
2739 .build(),
2740 )
2741 }
2742
2743 fn make_fill_voided(client_order_id: ClientOrderId, is_reopened: bool) -> OrderEventAny {
2744 OrderEventAny::FillVoided(
2745 OrderFillVoidedSpec::builder()
2746 .trader_id(TraderId::from("TRADER-001"))
2747 .strategy_id(StrategyId::from("TEST-001"))
2748 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2749 .client_order_id(client_order_id)
2750 .venue_order_id(VenueOrderId::test_default())
2751 .account_id(AccountId::from("ACC-001"))
2752 .is_reopened(is_reopened)
2753 .build(),
2754 )
2755 }
2756
2757 fn make_terminal_fill_voided(client_order_id: ClientOrderId) -> OrderEventAny {
2758 make_fill_voided(client_order_id, false)
2759 }
2760
2761 fn make_canceled(client_order_id: ClientOrderId) -> OrderEventAny {
2762 OrderEventAny::Canceled(
2763 OrderCanceledSpec::builder()
2764 .trader_id(TraderId::from("TRADER-001"))
2765 .strategy_id(StrategyId::from("TEST-001"))
2766 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2767 .client_order_id(client_order_id)
2768 .account_id(AccountId::from("ACC-001"))
2769 .event_id(UUID4::default())
2770 .build(),
2771 )
2772 }
2773
2774 fn make_rejected(client_order_id: ClientOrderId) -> OrderEventAny {
2775 OrderEventAny::Rejected(
2776 OrderRejectedSpec::builder()
2777 .trader_id(TraderId::from("TRADER-001"))
2778 .strategy_id(StrategyId::from("TEST-001"))
2779 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2780 .client_order_id(client_order_id)
2781 .account_id(AccountId::from("ACC-001"))
2782 .reason("Test rejection".into())
2783 .event_id(UUID4::default())
2784 .build(),
2785 )
2786 }
2787
2788 fn make_cancel_rejected(client_order_id: ClientOrderId) -> OrderEventAny {
2789 OrderEventAny::CancelRejected(
2790 OrderCancelRejectedSpec::builder()
2791 .trader_id(TraderId::from("TRADER-001"))
2792 .strategy_id(StrategyId::from("TEST-001"))
2793 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2794 .client_order_id(client_order_id)
2795 .reason("Test rejection".into())
2796 .venue_order_id(VenueOrderId::from(client_order_id.as_str()))
2797 .account_id(AccountId::from("ACC-001"))
2798 .build(),
2799 )
2800 }
2801
2802 fn make_expired(client_order_id: ClientOrderId) -> OrderEventAny {
2803 OrderEventAny::Expired(
2804 OrderExpiredSpec::builder()
2805 .trader_id(TraderId::from("TRADER-001"))
2806 .strategy_id(StrategyId::from("TEST-001"))
2807 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2808 .client_order_id(client_order_id)
2809 .account_id(AccountId::from("ACC-001"))
2810 .event_id(UUID4::default())
2811 .build(),
2812 )
2813 }
2814
2815 fn make_accepted(client_order_id: ClientOrderId) -> OrderEventAny {
2816 OrderEventAny::Accepted(
2817 OrderAcceptedSpec::builder()
2818 .trader_id(TraderId::from("TRADER-001"))
2819 .strategy_id(StrategyId::from("TEST-001"))
2820 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2821 .client_order_id(client_order_id)
2822 .venue_order_id(VenueOrderId::test_default())
2823 .account_id(AccountId::from("ACC-001"))
2824 .event_id(UUID4::default())
2825 .build(),
2826 )
2827 }
2828
2829 fn make_accepted_market_order(client_order_id: &str) -> OrderAny {
2830 let mut order = OrderAny::Market(MarketOrder::new(
2831 TraderId::from("TRADER-001"),
2832 StrategyId::from("TEST-001"),
2833 InstrumentId::from("BTCUSDT.BINANCE"),
2834 ClientOrderId::from(client_order_id),
2835 OrderSide::Buy,
2836 Quantity::from(100_000),
2837 TimeInForce::Gtc,
2838 UUID4::new(),
2839 UnixNanos::default(),
2840 false,
2841 false,
2842 None,
2843 None,
2844 None,
2845 None,
2846 None,
2847 None,
2848 None,
2849 None,
2850 ));
2851 let account_id = AccountId::from("ACC-001");
2852 order
2853 .apply(TestOrderEventStubs::submitted(&order, account_id))
2854 .unwrap();
2855 order
2856 .apply(TestOrderEventStubs::accepted(
2857 &order,
2858 account_id,
2859 VenueOrderId::from(client_order_id),
2862 ))
2863 .unwrap();
2864 order
2865 }
2866
2867 fn make_accepted_limit_order(client_order_id: &str) -> OrderAny {
2868 let mut order = OrderAny::Limit(LimitOrder::new(
2869 TraderId::from("TRADER-001"),
2870 StrategyId::from("TEST-001"),
2871 InstrumentId::from("BTCUSDT.BINANCE"),
2872 ClientOrderId::from(client_order_id),
2873 OrderSide::Buy,
2874 Quantity::from("1.0"),
2875 Price::from("50000.0"),
2876 TimeInForce::Gtc,
2877 None,
2878 false,
2879 false,
2880 false,
2881 None,
2882 None,
2883 None,
2884 None,
2885 None,
2886 None,
2887 None,
2888 None,
2889 None,
2890 None,
2891 None,
2892 UUID4::new(),
2893 UnixNanos::default(),
2894 ));
2895 let account_id = AccountId::from("ACC-001");
2896 order
2897 .apply(TestOrderEventStubs::submitted(&order, account_id))
2898 .unwrap();
2899 order
2900 .apply(TestOrderEventStubs::accepted(
2901 &order,
2902 account_id,
2903 VenueOrderId::from(client_order_id),
2905 ))
2906 .unwrap();
2907 order
2908 }
2909
2910 fn make_submitted_gtd_limit_order(client_order_id: &str, expire_time: UnixNanos) -> OrderAny {
2911 let mut order = OrderTestBuilder::new(OrderType::Limit)
2912 .trader_id(TraderId::from("TRADER-001"))
2913 .strategy_id(StrategyId::from("TEST-001"))
2914 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2915 .client_order_id(ClientOrderId::from(client_order_id))
2916 .side(OrderSide::Buy)
2917 .quantity(Quantity::from("1.0"))
2918 .price(Price::from("50000.0"))
2919 .time_in_force(TimeInForce::Gtd)
2920 .expire_time(expire_time)
2921 .build();
2922 let account_id = AccountId::from("ACC-001");
2923 order
2924 .apply(TestOrderEventStubs::submitted(&order, account_id))
2925 .unwrap();
2926 order
2927 }
2928
2929 fn make_accepted_gtd_limit_order(client_order_id: &str, expire_time: UnixNanos) -> OrderAny {
2930 let mut order = make_submitted_gtd_limit_order(client_order_id, expire_time);
2931 let account_id = AccountId::from("ACC-001");
2932 order
2933 .apply(TestOrderEventStubs::accepted(
2934 &order,
2935 account_id,
2936 VenueOrderId::from(client_order_id),
2937 ))
2938 .unwrap();
2939 order
2940 }
2941
2942 fn make_initialized_market_order(client_order_id: &str) -> OrderAny {
2943 OrderAny::Market(MarketOrder::new(
2944 TraderId::from("TRADER-001"),
2945 StrategyId::from("TEST-001"),
2946 InstrumentId::from("BTCUSDT.BINANCE"),
2947 ClientOrderId::from(client_order_id),
2948 OrderSide::Buy,
2949 Quantity::from(100_000),
2950 TimeInForce::Gtc,
2951 UUID4::new(),
2952 UnixNanos::default(),
2953 false,
2954 false,
2955 None,
2956 None,
2957 None,
2958 None,
2959 None,
2960 None,
2961 None,
2962 None,
2963 ))
2964 }
2965
2966 fn make_initialized_algorithm_order(client_order_id: &str) -> OrderAny {
2967 OrderAny::Market(MarketOrder::new(
2968 TraderId::from("TRADER-001"),
2969 StrategyId::from("TEST-001"),
2970 InstrumentId::from("BTCUSDT.BINANCE"),
2971 ClientOrderId::from(client_order_id),
2972 OrderSide::Buy,
2973 Quantity::from(100_000),
2974 TimeInForce::Gtc,
2975 UUID4::new(),
2976 UnixNanos::default(),
2977 false,
2978 false,
2979 None,
2980 None,
2981 None,
2982 None,
2983 Some(ExecAlgorithmId::from("TWAP")),
2984 None,
2985 Some(ClientOrderId::from(client_order_id)),
2986 None,
2987 ))
2988 }
2989
2990 fn add_order_to_cache(strategy: &TestStrategy, order: &OrderAny) {
2991 let cache_rc = strategy.core.cache_rc();
2992 let mut cache = cache_rc.borrow_mut();
2993 cache.add_order(order.clone(), None, None, true).unwrap();
2994 }
2995
2996 fn make_submit_command(order: &OrderAny) -> SubmitOrder {
2997 SubmitOrder::new(
2998 order.trader_id(),
2999 None,
3000 order.strategy_id(),
3001 order.instrument_id(),
3002 order.client_order_id(),
3003 order.init_event().clone(),
3004 order.exec_algorithm_id(),
3005 None,
3006 None,
3007 UUID4::new(),
3008 UnixNanos::default(),
3009 None, )
3011 }
3012
3013 fn add_order_to_cache_and_own_book(strategy: &TestStrategy, order: &OrderAny) {
3014 let cache_rc = strategy.core.cache_rc();
3015 let mut cache = cache_rc.borrow_mut();
3016 cache.add_order(order.clone(), None, None, true).unwrap();
3017 cache
3018 .add_own_order_book(OwnOrderBook::new(order.instrument_id()))
3019 .unwrap();
3020 cache.update_own_order_book(order);
3021 }
3022
3023 fn make_position_opened() -> PositionEvent {
3024 PositionEvent::PositionOpened(PositionOpened {
3025 trader_id: TraderId::from("TRADER-001"),
3026 strategy_id: StrategyId::from("TEST-001"),
3027 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
3028 position_id: PositionId::test_default(),
3029 account_id: AccountId::from("ACC-001"),
3030 opening_order_id: ClientOrderId::from("O-001"),
3031 entry: OrderSide::Buy,
3032 side: PositionSide::Long,
3033 signed_qty: 1.0,
3034 quantity: Quantity::default(),
3035 last_qty: Quantity::default(),
3036 last_px: Price::default(),
3037 currency: Currency::from("USD"),
3038 avg_px_open: 0.0,
3039 realized_pnl: None,
3040 event_id: UUID4::default(),
3041 ts_event: UnixNanos::default(),
3042 ts_init: UnixNanos::default(),
3043 })
3044 }
3045
3046 fn make_position_changed() -> PositionEvent {
3047 let currency = Currency::from("USD");
3048 PositionEvent::PositionChanged(PositionChanged {
3049 trader_id: TraderId::from("TRADER-001"),
3050 strategy_id: StrategyId::from("TEST-001"),
3051 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
3052 position_id: PositionId::test_default(),
3053 account_id: AccountId::from("ACC-001"),
3054 opening_order_id: ClientOrderId::from("O-001"),
3055 entry: OrderSide::Buy,
3056 side: PositionSide::Long,
3057 signed_qty: 2.0,
3058 quantity: Quantity::default(),
3059 peak_quantity: Quantity::default(),
3060 last_qty: Quantity::default(),
3061 last_px: Price::default(),
3062 currency,
3063 avg_px_open: 0.0,
3064 avg_px_close: None,
3065 realized_return: 0.0,
3066 realized_pnl: None,
3067 unrealized_pnl: Money::zero(currency),
3068 event_id: UUID4::default(),
3069 ts_opened: UnixNanos::default(),
3070 ts_event: UnixNanos::default(),
3071 ts_init: UnixNanos::default(),
3072 })
3073 }
3074
3075 fn make_position_closed() -> PositionEvent {
3076 let currency = Currency::from("USD");
3077 PositionEvent::PositionClosed(PositionClosed {
3078 trader_id: TraderId::from("TRADER-001"),
3079 strategy_id: StrategyId::from("TEST-001"),
3080 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
3081 position_id: PositionId::test_default(),
3082 account_id: AccountId::from("ACC-001"),
3083 opening_order_id: ClientOrderId::from("O-001"),
3084 closing_order_id: Some(ClientOrderId::from("O-002")),
3085 entry: OrderSide::Buy,
3086 side: PositionSide::Flat,
3087 signed_qty: 0.0,
3088 quantity: Quantity::default(),
3089 peak_quantity: Quantity::default(),
3090 last_qty: Quantity::default(),
3091 last_px: Price::default(),
3092 currency,
3093 avg_px_open: 0.0,
3094 avg_px_close: None,
3095 realized_return: 0.0,
3096 realized_pnl: None,
3097 unrealized_pnl: Money::zero(currency),
3098 duration: DurationNanos::default(),
3099 event_id: UUID4::default(),
3100 ts_opened: UnixNanos::default(),
3101 ts_closed: None,
3102 ts_event: UnixNanos::default(),
3103 ts_init: UnixNanos::default(),
3104 })
3105 }
3106
3107 fn make_position_adjusted() -> PositionEvent {
3108 PositionEvent::PositionAdjusted(PositionAdjusted {
3109 trader_id: TraderId::from("TRADER-001"),
3110 strategy_id: StrategyId::from("TEST-001"),
3111 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
3112 position_id: PositionId::test_default(),
3113 account_id: AccountId::from("ACC-001"),
3114 adjustment_type: PositionAdjustmentType::Funding,
3115 quantity_change: None,
3116 pnl_change: None,
3117 reason: None,
3118 event_id: UUID4::default(),
3119 ts_event: UnixNanos::default(),
3120 ts_init: UnixNanos::default(),
3121 })
3122 }
3123
3124 #[rstest]
3125 fn test_strategy_creation() {
3126 let strategy = create_test_strategy();
3127 assert_eq!(strategy.strategy_id(), Some(StrategyId::from("TEST-001")));
3128 assert!(!strategy.on_order_rejected_called);
3129 assert!(!strategy.on_position_opened_called);
3130 }
3131
3132 #[rstest]
3133 fn test_strategy_registration() {
3134 let mut strategy = create_test_strategy();
3135 register_strategy(&mut strategy);
3136
3137 assert!(strategy.is_registered());
3138 let _ = strategy.order().generate_client_order_id();
3139 let _ = strategy.portfolio().is_initialized();
3140 }
3141
3142 #[rstest]
3143 fn test_set_external_order_instrument_ids_replaces_claims_atomically() {
3144 let mut strategy = create_test_strategy();
3145 register_strategy(&mut strategy);
3146 let cache = strategy.core.cache_rc();
3147 let strategy_id = StrategyId::from("TEST-001");
3148 let other_strategy_id = StrategyId::from("OTHER-001");
3149 let audusd = InstrumentId::from("AUDUSD.SIM");
3150 let eurusd = InstrumentId::from("EURUSD.SIM");
3151 let gbpusd = InstrumentId::from("GBPUSD.SIM");
3152
3153 strategy
3154 .set_external_order_instrument_ids(vec![audusd, eurusd])
3155 .unwrap();
3156 cache
3157 .borrow_mut()
3158 .set_external_order_claims(other_strategy_id, &[gbpusd])
3159 .unwrap();
3160
3161 let result = strategy.set_external_order_instrument_ids(vec![eurusd, gbpusd]);
3162
3163 assert_eq!(
3164 result.unwrap_err().to_string(),
3165 "External order claim for GBPUSD.SIM already exists for OTHER-001"
3166 );
3167 assert_eq!(
3168 cache.borrow().external_order_claim(&audusd),
3169 Some(strategy_id)
3170 );
3171 assert_eq!(
3172 cache.borrow().external_order_claim(&eurusd),
3173 Some(strategy_id)
3174 );
3175 assert_eq!(
3176 cache.borrow().external_order_claim(&gbpusd),
3177 Some(other_strategy_id)
3178 );
3179 assert_eq!(
3180 strategy.core.config.external_order_instrument_ids,
3181 Some(vec![audusd, eurusd])
3182 );
3183
3184 strategy
3185 .set_external_order_instrument_ids(vec![eurusd])
3186 .unwrap();
3187
3188 assert_eq!(cache.borrow().external_order_claim(&audusd), None);
3189 assert_eq!(
3190 cache.borrow().external_order_claim(&eurusd),
3191 Some(strategy_id)
3192 );
3193 assert_eq!(
3194 cache.borrow().external_order_claim(&gbpusd),
3195 Some(other_strategy_id)
3196 );
3197 assert_eq!(
3198 strategy.core.config.external_order_instrument_ids,
3199 Some(vec![eurusd])
3200 );
3201 }
3202
3203 #[rstest]
3204 fn test_set_external_order_instrument_ids_rejects_unregistered_strategy() {
3205 let mut strategy = create_test_strategy();
3206
3207 let error = strategy
3208 .set_external_order_instrument_ids(vec![InstrumentId::from("AUDUSD.SIM")])
3209 .unwrap_err();
3210
3211 assert_eq!(
3212 error.to_string(),
3213 "Strategy TEST-001 is not registered with a trader"
3214 );
3215 assert!(strategy.core.config.external_order_instrument_ids.is_none());
3216 }
3217
3218 #[rstest]
3219 fn test_strategy_native_methods_are_available_on_strategy_type() {
3220 let mut strategy = create_test_strategy();
3221 register_strategy(&mut strategy);
3222
3223 drop(strategy.order_factory());
3224
3225 assert!(Rc::ptr_eq(
3226 &strategy.order_factory_rc(),
3227 strategy.core.order_factory.as_ref().unwrap()
3228 ));
3229 assert!(Rc::ptr_eq(
3230 &strategy.portfolio_rc(),
3231 strategy.core.portfolio.as_ref().unwrap()
3232 ));
3233 }
3234
3235 #[rstest]
3236 fn test_handle_order_event_dispatches_to_handler() {
3237 let mut strategy = create_test_strategy();
3238 register_strategy(&mut strategy);
3239 start_strategy(&mut strategy);
3240
3241 let event = make_rejected(ClientOrderId::from("O-001"));
3242
3243 strategy.handle_order_event(event);
3244
3245 assert!(strategy.on_order_rejected_called);
3246 assert!(strategy.on_order_event_called);
3247 }
3248
3249 #[rstest]
3250 fn test_dispatch_manager_actions_routes_every_action() {
3251 let mut strategy = create_test_strategy();
3252 register_strategy(&mut strategy);
3253 let strategy_id = StrategyId::from("TEST-001");
3254 let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
3255 let initialized_order = make_initialized_market_order("O-MANAGER-INIT");
3256 let emulator_order = OrderTestBuilder::new(OrderType::StopMarket)
3257 .trader_id(TraderId::from("TRADER-001"))
3258 .strategy_id(strategy_id)
3259 .instrument_id(instrument_id)
3260 .client_order_id(ClientOrderId::from("O-MANAGER-EMULATOR"))
3261 .side(OrderSide::Buy)
3262 .trigger_price(Price::from("51000.0"))
3263 .quantity(Quantity::from(100_000))
3264 .emulation_trigger(TriggerType::BidAsk)
3265 .build();
3266 let risk_order = make_initialized_market_order("O-MANAGER-RISK");
3267 let algorithm_order = make_initialized_algorithm_order("O-MANAGER-ALGORITHM");
3268 let cancel_order = OrderTestBuilder::new(OrderType::Limit)
3269 .trader_id(TraderId::from("TRADER-001"))
3270 .strategy_id(strategy_id)
3271 .instrument_id(instrument_id)
3272 .client_order_id(ClientOrderId::from("O-MANAGER-CANCEL"))
3273 .side(OrderSide::Buy)
3274 .price(Price::from("50000.0"))
3275 .quantity(Quantity::from(100_000))
3276 .submit(true)
3277 .build();
3278 let missing_cancel_order = OrderTestBuilder::new(OrderType::Limit)
3279 .trader_id(TraderId::from("TRADER-001"))
3280 .strategy_id(strategy_id)
3281 .instrument_id(instrument_id)
3282 .client_order_id(ClientOrderId::from("O-MANAGER-MISSING-CANCEL"))
3283 .side(OrderSide::Buy)
3284 .price(Price::from("50000.0"))
3285 .quantity(Quantity::from(100_000))
3286 .submit(true)
3287 .build();
3288 let modify_order = OrderTestBuilder::new(OrderType::Limit)
3289 .trader_id(TraderId::from("TRADER-001"))
3290 .strategy_id(strategy_id)
3291 .instrument_id(instrument_id)
3292 .client_order_id(ClientOrderId::from("O-MANAGER-MODIFY"))
3293 .side(OrderSide::Buy)
3294 .price(Price::from("50000.0"))
3295 .quantity(Quantity::from(100_000))
3296 .submit(true)
3297 .build();
3298 add_order_to_cache(&strategy, &cancel_order);
3299 add_order_to_cache(&strategy, &modify_order);
3300
3301 let (order_handler, order_events) =
3302 get_typed_message_saving_handler(Some(Ustr::from("manager-action-order-events")));
3303 let order_topic = format!("events.order.{strategy_id}");
3304 msgbus::subscribe_order_events(order_topic.clone().into(), order_handler.clone(), None);
3305 let (emulator_handler, emulator_messages): (
3306 _,
3307 TypedIntoMessageSavingHandler<TradingCommand>,
3308 ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
3309 msgbus::register_trading_command_endpoint(
3310 MessagingSwitchboard::order_emulator_execute(),
3311 emulator_handler,
3312 );
3313 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3314 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3315 msgbus::register_trading_command_endpoint(
3316 MessagingSwitchboard::risk_engine_queue_execute(),
3317 risk_handler,
3318 );
3319 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3320 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3321 msgbus::register_trading_command_endpoint(
3322 MessagingSwitchboard::exec_engine_queue_execute(),
3323 exec_handler,
3324 );
3325 let (algorithm_handler, algorithm_messages) =
3326 get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
3327 msgbus::register_any("TWAP.execute".into(), algorithm_handler);
3328
3329 strategy.dispatch_manager_actions(vec![
3330 OrderManagerAction::CancelLocal(missing_cancel_order),
3331 OrderManagerAction::PublishInitialized(OrderEventAny::Initialized(
3332 initialized_order.init_event().clone(),
3333 )),
3334 OrderManagerAction::SubmitToEmulator(make_submit_command(&emulator_order)),
3335 OrderManagerAction::SubmitToRisk(make_submit_command(&risk_order)),
3336 OrderManagerAction::SubmitToAlgorithm {
3337 command: make_submit_command(&algorithm_order),
3338 exec_algorithm_id: ExecAlgorithmId::from("TWAP"),
3339 },
3340 OrderManagerAction::CancelLocal(cancel_order.clone()),
3341 OrderManagerAction::ModifyLocalQuantity {
3342 order: modify_order.clone(),
3343 quantity: Quantity::from(50_000),
3344 },
3345 OrderManagerAction::ModifyLocalQuantity {
3346 order: modify_order.clone(),
3347 quantity: Quantity::from(100_000),
3348 },
3349 ]);
3350 msgbus::unsubscribe_order_events(order_topic.into(), &order_handler);
3351
3352 let order_events = order_events.get_messages();
3353 assert_eq!(order_events.len(), 3);
3354 assert!(matches!(
3355 &order_events[0],
3356 OrderEventAny::Initialized(event)
3357 if event.client_order_id == initialized_order.client_order_id()
3358 ));
3359 assert!(matches!(
3360 &order_events[1],
3361 OrderEventAny::PendingCancel(event)
3362 if event.client_order_id == cancel_order.client_order_id()
3363 ));
3364 assert!(matches!(
3365 &order_events[2],
3366 OrderEventAny::PendingUpdate(event)
3367 if event.client_order_id == modify_order.client_order_id()
3368 ));
3369 assert!(matches!(
3370 emulator_messages.get_messages().as_slice(),
3371 [TradingCommand::SubmitOrder(command)]
3372 if command.client_order_id == emulator_order.client_order_id()
3373 ));
3374 assert!(matches!(
3375 risk_messages.get_messages().as_slice(),
3376 [
3377 TradingCommand::SubmitOrder(submit),
3378 TradingCommand::ModifyOrder(first_modify),
3379 TradingCommand::ModifyOrder(second_modify),
3380 ]
3381
3382 if submit.client_order_id == risk_order.client_order_id()
3383 && first_modify.client_order_id == modify_order.client_order_id()
3384 && first_modify.quantity == Some(Quantity::from(50_000))
3385 && second_modify.client_order_id == modify_order.client_order_id()
3386 && second_modify.quantity == Some(Quantity::from(100_000))
3387 ));
3388 assert!(matches!(
3389 algorithm_messages.get_messages().as_slice(),
3390 [TradingCommand::SubmitOrder(command)]
3391 if command.client_order_id == algorithm_order.client_order_id()
3392 ));
3393 assert!(matches!(
3394 exec_messages.get_messages().as_slice(),
3395 [TradingCommand::CancelOrder(command)]
3396 if command.client_order_id == cancel_order.client_order_id()
3397 ));
3398 }
3399
3400 #[rstest]
3401 #[case::disabled(false)]
3402 #[case::enabled(true)]
3403 fn test_contingent_manager_dispatch_precedes_user_handlers_and_is_idempotent(
3404 #[case] manage_contingent_orders: bool,
3405 ) {
3406 let mut strategy = TestStrategy::new(StrategyConfig {
3407 strategy_id: Some(StrategyId::from("TEST-001")),
3408 order_id_tag: Some("001".to_string()),
3409 manage_contingent_orders,
3410 ..Default::default()
3411 });
3412 register_strategy(&mut strategy);
3413 start_strategy(&mut strategy);
3414 let timeline = strategy.order_event_timeline.clone();
3415 let endpoint_timeline = timeline.clone();
3416 let exec_handler = TypedIntoHandler::from(move |command: TradingCommand| {
3417 assert!(matches!(command, TradingCommand::CancelOrder(_)));
3418 endpoint_timeline.borrow_mut().push("endpoint");
3419 });
3420 msgbus::register_trading_command_endpoint(
3421 MessagingSwitchboard::exec_engine_queue_execute(),
3422 exec_handler,
3423 );
3424 let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
3425 let parent_id = ClientOrderId::from("O-CONTINGENT-PARENT");
3426 let sibling_id = ClientOrderId::from("O-CONTINGENT-SIBLING");
3427 let parent = OrderTestBuilder::new(OrderType::Limit)
3428 .trader_id(TraderId::from("TRADER-001"))
3429 .strategy_id(StrategyId::from("TEST-001"))
3430 .instrument_id(instrument_id)
3431 .client_order_id(parent_id)
3432 .side(OrderSide::Buy)
3433 .price(Price::from("50000.0"))
3434 .quantity(Quantity::from(100_000))
3435 .contingency_type(ContingencyType::Oco)
3436 .linked_order_ids(vec![parent_id, sibling_id])
3437 .submit(true)
3438 .build();
3439 let sibling = OrderTestBuilder::new(OrderType::Limit)
3440 .trader_id(TraderId::from("TRADER-001"))
3441 .strategy_id(StrategyId::from("TEST-001"))
3442 .instrument_id(instrument_id)
3443 .client_order_id(sibling_id)
3444 .side(OrderSide::Sell)
3445 .price(Price::from("51000.0"))
3446 .quantity(Quantity::from(100_000))
3447 .submit(true)
3448 .build();
3449 add_order_to_cache(&strategy, &parent);
3450 add_order_to_cache(&strategy, &sibling);
3451 let event = OrderEventAny::Filled(
3452 OrderFilledSpec::builder()
3453 .trader_id(parent.trader_id())
3454 .strategy_id(parent.strategy_id())
3455 .instrument_id(parent.instrument_id())
3456 .client_order_id(parent.client_order_id())
3457 .venue_order_id(VenueOrderId::from("V-CONTINGENT-PARENT"))
3458 .account_id(AccountId::from("ACCOUNT-001"))
3459 .trade_id(TradeId::from("T-CONTINGENT-PARENT"))
3460 .order_side(parent.order_side())
3461 .order_type(parent.order_type())
3462 .last_qty(parent.quantity())
3463 .last_px(Price::from("50000.0"))
3464 .liquidity_side(LiquiditySide::Taker)
3465 .build(),
3466 );
3467 strategy
3468 .core
3469 .cache_rc()
3470 .borrow_mut()
3471 .update_order(&event)
3472 .unwrap();
3473
3474 strategy.handle_order_event(event.clone());
3475 strategy.handle_order_event(event);
3476
3477 let timeline = timeline.borrow();
3478 let sibling_status = strategy
3479 .core
3480 .cache_ref()
3481 .order(&sibling_id)
3482 .unwrap()
3483 .status();
3484
3485 if manage_contingent_orders {
3486 assert_eq!(
3487 timeline.as_slice(),
3488 ["endpoint", "specific", "aggregate", "specific", "aggregate"]
3489 );
3490 assert_eq!(sibling_status, OrderStatus::PendingCancel);
3491 } else {
3492 assert_eq!(
3493 timeline.as_slice(),
3494 ["specific", "aggregate", "specific", "aggregate"]
3495 );
3496 assert_eq!(sibling_status, OrderStatus::Submitted);
3497 }
3498 }
3499
3500 #[rstest]
3501 fn test_handle_order_fill_voided_dispatches_to_specific_handler() {
3502 let mut strategy = create_test_strategy();
3503 register_strategy(&mut strategy);
3504 start_strategy(&mut strategy);
3505
3506 strategy.handle_order_event(make_fill_voided(ClientOrderId::from("O-001"), false));
3507
3508 assert!(strategy.on_order_fill_voided_called);
3509 assert!(strategy.on_order_event_called);
3510 }
3511
3512 #[rstest]
3513 #[case::opened(make_position_opened())]
3514 #[case::changed(make_position_changed())]
3515 #[case::closed(make_position_closed())]
3516 fn test_handle_position_event_dispatches_to_handler(#[case] event: PositionEvent) {
3517 let mut strategy = create_test_strategy();
3518 register_strategy(&mut strategy);
3519 start_strategy(&mut strategy);
3520
3521 let expected_opened = matches!(event, PositionEvent::PositionOpened(_));
3522 let expected_changed = matches!(event, PositionEvent::PositionChanged(_));
3523 let expected_closed = matches!(event, PositionEvent::PositionClosed(_));
3524
3525 strategy.handle_position_event(event);
3526
3527 assert_eq!(strategy.on_position_opened_called, expected_opened);
3528 assert_eq!(strategy.on_position_changed_called, expected_changed);
3529 assert_eq!(strategy.on_position_closed_called, expected_closed);
3530 assert!(strategy.on_position_event_called);
3531 }
3532
3533 #[rstest]
3534 fn test_handle_position_event_skips_dispatch_when_stopped() {
3535 let mut strategy = create_test_strategy();
3536 register_strategy(&mut strategy);
3537 start_strategy(&mut strategy);
3538 stop_strategy(&mut strategy);
3539 assert_eq!(strategy.state(), ComponentState::Stopped);
3540
3541 strategy.handle_position_event(make_position_opened());
3542
3543 assert!(!strategy.on_position_event_called);
3544 assert!(!strategy.on_position_opened_called);
3545 }
3546
3547 #[rstest]
3548 fn test_handle_position_event_skips_dispatch_for_adjusted() {
3549 let mut strategy = create_test_strategy();
3550 register_strategy(&mut strategy);
3551 start_strategy(&mut strategy);
3552
3553 strategy.handle_position_event(make_position_adjusted());
3554
3555 assert!(!strategy.on_position_event_called);
3556 assert!(!strategy.on_position_opened_called);
3557 assert!(!strategy.on_position_changed_called);
3558 assert!(!strategy.on_position_closed_called);
3559 }
3560
3561 #[rstest]
3562 fn test_strategy_default_handlers_do_not_panic() {
3563 let mut strategy = create_test_strategy();
3564
3565 strategy.on_order_initialized(OrderInitialized::default());
3566 strategy.on_order_event(OrderEventAny::Accepted(OrderAccepted::default()));
3567 strategy.on_order_denied(OrderDenied::default());
3568 strategy.on_order_emulated(OrderEmulated::default());
3569 strategy.on_order_released(OrderReleased::default());
3570 strategy.on_order_submitted(OrderSubmitted::default());
3571 strategy.on_order_rejected(OrderRejected::default());
3572 strategy.on_order_canceled(&OrderCanceled::default());
3573 strategy.on_order_expired(OrderExpired::default());
3574 strategy.on_order_triggered(OrderTriggered::default());
3575 strategy.on_order_pending_update(OrderPendingUpdate::default());
3576 strategy.on_order_pending_cancel(OrderPendingCancel::default());
3577 strategy.on_order_modify_rejected(OrderModifyRejected::default());
3578 strategy.on_order_cancel_rejected(OrderCancelRejected::default());
3579 strategy.on_order_updated(OrderUpdated::default());
3580 strategy.on_order_filled(&OrderFilledSpec::builder().build());
3581 strategy.on_order_fill_voided(&OrderFillVoidedSpec::builder().build());
3582 strategy.on_position_event(make_position_opened());
3583 }
3584
3585 #[rstest]
3586 fn test_submit_order_publishes_order_initialized_after_cache_insert_before_send() {
3587 let mut strategy = create_test_strategy();
3588 register_strategy(&mut strategy);
3589
3590 let order = make_initialized_market_order("O-20250208-INIT-001");
3591 let client_order_id = order.client_order_id();
3592 let cache_rc = strategy.core.cache_rc();
3593 let timeline = Rc::new(RefCell::new(Vec::new()));
3594 let event_messages = Rc::new(RefCell::new(Vec::new()));
3595
3596 let event_handler = {
3597 let event_messages = event_messages.clone();
3598 let timeline = timeline.clone();
3599 TypedHandler::from_with_id("events.order.initialized", move |event: &OrderEventAny| {
3600 assert!(cache_rc.borrow().order_exists(&client_order_id));
3601 assert!(matches!(event, OrderEventAny::Initialized(_)));
3602 event_messages.borrow_mut().push(event.clone());
3603 timeline.borrow_mut().push("init");
3604 })
3605 };
3606 let risk_handler = {
3607 let timeline = timeline.clone();
3608 TypedIntoHandler::from_with_id(
3609 "RiskEngine.queue_execute",
3610 move |command: TradingCommand| {
3611 assert!(matches!(command, TradingCommand::SubmitOrder(_)));
3612 timeline.borrow_mut().push("command");
3613 },
3614 )
3615 };
3616 msgbus::register_trading_command_endpoint(
3617 MessagingSwitchboard::risk_engine_queue_execute(),
3618 risk_handler,
3619 );
3620
3621 let topic = format!("events.order.{}", order.strategy_id());
3622 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3623
3624 strategy
3625 .submit_order(order.clone(), None, None, None)
3626 .unwrap();
3627
3628 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3629
3630 let event_messages = event_messages.borrow();
3631 assert_eq!(event_messages.len(), 1);
3632 assert_eq!(
3633 event_messages[0],
3634 OrderEventAny::Initialized(order.init_event().clone())
3635 );
3636 assert_eq!(timeline.borrow().as_slice(), &["init", "command"]);
3637 }
3638
3639 #[rstest]
3640 fn test_submit_order_routes_emulated_order_to_order_emulator() {
3641 let mut strategy = create_test_strategy();
3642 register_strategy(&mut strategy);
3643 let (emulator_handler, emulator_messages): (
3644 _,
3645 TypedIntoMessageSavingHandler<TradingCommand>,
3646 ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
3647 msgbus::register_trading_command_endpoint(
3648 MessagingSwitchboard::order_emulator_execute(),
3649 emulator_handler,
3650 );
3651 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3652 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3653 msgbus::register_trading_command_endpoint(
3654 MessagingSwitchboard::risk_engine_queue_execute(),
3655 risk_handler,
3656 );
3657 let order = OrderTestBuilder::new(OrderType::StopMarket)
3658 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
3659 .client_order_id(ClientOrderId::from("O-20250208-EMULATED-001"))
3660 .side(OrderSide::Buy)
3661 .trigger_price(Price::from("51000.0"))
3662 .quantity(Quantity::from(100_000))
3663 .emulation_trigger(TriggerType::BidAsk)
3664 .build();
3665 let client_order_id = order.client_order_id();
3666
3667 strategy.submit_order(order, None, None, None).unwrap();
3668
3669 let emulator_messages = emulator_messages.get_messages();
3670 assert_eq!(emulator_messages.len(), 1);
3671 assert!(matches!(
3672 emulator_messages.first(),
3673 Some(TradingCommand::SubmitOrder(command))
3674 if command.client_order_id == client_order_id
3675 ));
3676 assert!(risk_messages.get_messages().is_empty());
3677 }
3678
3679 #[rstest]
3680 fn test_submit_order_errors_when_strategy_not_registered() {
3681 let mut strategy = create_test_strategy();
3682 let order = make_initialized_market_order("O-20250208-UNREGISTERED-001");
3683
3684 let err = strategy
3685 .submit_order(order, None, None, None)
3686 .unwrap_err()
3687 .to_string();
3688
3689 assert_eq!(err, "Strategy not registered: trader_id is not set");
3690 }
3691
3692 #[rstest]
3693 fn test_submit_order_uses_stored_strategy_id_when_actor_id_diverges() {
3694 let mut strategy = create_test_strategy();
3695 register_strategy(&mut strategy);
3696 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3697 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3698 msgbus::register_trading_command_endpoint(
3699 MessagingSwitchboard::risk_engine_queue_execute(),
3700 risk_handler,
3701 );
3702
3703 strategy.core.actor.actor_id = ActorId::from("Strategy");
3705 let order = make_initialized_market_order("O-20250208-DIVERGED-001");
3706
3707 strategy.submit_order(order, None, None, None).unwrap();
3708
3709 let risk_messages = risk_messages.get_messages();
3710 assert_eq!(risk_messages.len(), 1);
3711 let Some(TradingCommand::SubmitOrder(command)) = risk_messages.first() else {
3712 panic!("Expected a SubmitOrder command, was {risk_messages:?}");
3713 };
3714 assert_eq!(command.strategy_id, StrategyId::from("TEST-001"));
3715 assert_eq!(
3716 command.client_order_id,
3717 ClientOrderId::from("O-20250208-DIVERGED-001")
3718 );
3719 }
3720
3721 #[rstest]
3722 fn test_required_account_id_errors_when_missing_for_strategy_event() {
3723 let order = make_initialized_market_order("O-20250208-NO-ACCOUNT-001");
3724
3725 let err = required_account_id(&order, "pending cancel")
3726 .unwrap_err()
3727 .to_string();
3728
3729 assert_eq!(
3730 err,
3731 "Cannot generate pending cancel event for O-20250208-NO-ACCOUNT-001: \
3732 account_id is not set"
3733 );
3734 }
3735
3736 #[rstest]
3737 fn test_submit_order_rejects_non_initialized_without_events() {
3738 let mut strategy = create_test_strategy();
3739 register_strategy(&mut strategy);
3740
3741 let order = make_accepted_market_order("O-20250208-ACCEPTED-001");
3742 let topic = format!("events.order.{}", order.strategy_id());
3743 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
3744 get_typed_message_saving_handler(Some(Ustr::from("events.order.invalid")));
3745
3746 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3747 let result = strategy.submit_order(order, None, None, None);
3748
3749 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3750
3751 assert!(result.is_err());
3752 assert!(
3753 result
3754 .unwrap_err()
3755 .to_string()
3756 .contains("expected INITIALIZED")
3757 );
3758 assert!(event_messages.get_messages().is_empty());
3759 }
3760
3761 #[rstest]
3762 fn test_submit_order_returns_error_when_cache_already_borrowed() {
3763 let mut strategy = create_test_strategy();
3764 register_strategy(&mut strategy);
3765
3766 let order = make_initialized_market_order("O-20250208-BORROWED-001");
3767 let cache_rc = strategy.core.cache_rc();
3768 let _cache = cache_rc.borrow();
3769
3770 let result = catch_unwind(AssertUnwindSafe(|| {
3771 strategy.submit_order(order, None, None, None)
3772 }));
3773
3774 let err = result
3775 .expect("submit_order should not panic")
3776 .unwrap_err()
3777 .to_string();
3778
3779 assert_eq!(
3780 err,
3781 "Cannot submit order O-20250208-BORROWED-001: cache is currently borrowed"
3782 );
3783 }
3784
3785 #[rstest]
3786 fn test_submit_order_list_publishes_order_initialized_after_cache_insert_before_send() {
3787 let mut strategy = create_test_strategy();
3788 register_strategy(&mut strategy);
3789
3790 let order_list_id = OrderListId::from("OL-20250208-LIST-INIT");
3791 let mut orders = vec![
3792 make_initialized_market_order("O-20250208-LIST-INIT-001"),
3793 make_initialized_market_order("O-20250208-LIST-INIT-002"),
3794 ];
3795
3796 for order in &mut orders {
3797 order.set_order_list_id(order_list_id);
3798 }
3799
3800 let client_order_id1 = orders[0].client_order_id();
3801 let client_order_id2 = orders[1].client_order_id();
3802 let cache_rc = strategy.core.cache_rc();
3803 let timeline = Rc::new(RefCell::new(Vec::new()));
3804 let event_messages = Rc::new(RefCell::new(Vec::new()));
3805
3806 let event_handler = {
3807 let event_messages = event_messages.clone();
3808 let timeline = timeline.clone();
3809 TypedHandler::from_with_id(
3810 "events.order.list_initialized",
3811 move |event: &OrderEventAny| {
3812 match event {
3813 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id1 => {
3814 let cache = cache_rc.borrow();
3815 assert!(cache.order_exists(&client_order_id1));
3816 assert!(cache.order_exists(&client_order_id2));
3817 assert!(cache.order_list_exists(&order_list_id));
3818 let order_list = cache.order_list(&order_list_id).unwrap();
3819 assert_eq!(
3820 order_list.client_order_ids.as_slice(),
3821 &[client_order_id1, client_order_id2]
3822 );
3823 timeline.borrow_mut().push("init1");
3824 }
3825 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id2 => {
3826 assert!(cache_rc.borrow().order_exists(&client_order_id2));
3827 timeline.borrow_mut().push("init2");
3828 }
3829 _ => panic!("unexpected order event {event:?}"),
3830 }
3831 event_messages.borrow_mut().push(event.clone());
3832 },
3833 )
3834 };
3835 let risk_handler = {
3836 let timeline = timeline.clone();
3837 TypedIntoHandler::from_with_id(
3838 "RiskEngine.queue_execute",
3839 move |command: TradingCommand| {
3840 assert!(matches!(command, TradingCommand::SubmitOrderList(_)));
3841 timeline.borrow_mut().push("command");
3842 },
3843 )
3844 };
3845 msgbus::register_trading_command_endpoint(
3846 MessagingSwitchboard::risk_engine_queue_execute(),
3847 risk_handler,
3848 );
3849
3850 let topic = format!("events.order.{}", orders[0].strategy_id());
3851 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3852
3853 strategy
3854 .submit_order_list(orders.clone(), None, None, None)
3855 .unwrap();
3856
3857 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3858
3859 let event_messages = event_messages.borrow();
3860 assert_eq!(event_messages.len(), 2);
3861 assert_eq!(
3862 event_messages[0],
3863 OrderEventAny::Initialized(orders[0].init_event().clone())
3864 );
3865 assert_eq!(
3866 event_messages[1],
3867 OrderEventAny::Initialized(orders[1].init_event().clone())
3868 );
3869 assert_eq!(timeline.borrow().as_slice(), &["init1", "init2", "command"]);
3870 }
3871
3872 #[rstest]
3873 fn test_submit_order_list_returns_error_when_cache_already_borrowed() {
3874 let mut strategy = create_test_strategy();
3875 register_strategy(&mut strategy);
3876
3877 let order_list_id = OrderListId::from("OL-20250208-BORROWED");
3878 let mut orders = vec![
3879 make_initialized_market_order("O-20250208-LIST-BORROWED-001"),
3880 make_initialized_market_order("O-20250208-LIST-BORROWED-002"),
3881 ];
3882
3883 for order in &mut orders {
3884 order.set_order_list_id(order_list_id);
3885 }
3886
3887 let cache_rc = strategy.core.cache_rc();
3888 let _cache = cache_rc.borrow();
3889
3890 let result = catch_unwind(AssertUnwindSafe(|| {
3891 strategy.submit_order_list(orders, None, None, None)
3892 }));
3893
3894 let err = result
3895 .expect("submit_order_list should not panic")
3896 .unwrap_err()
3897 .to_string();
3898
3899 assert_eq!(
3900 err,
3901 "Cannot submit order list OL-20250208-BORROWED: cache is currently borrowed"
3902 );
3903 }
3904
3905 #[rstest]
3906 fn test_submit_order_list_create_list_branch_publishes_init_after_cache_insert() {
3907 let mut strategy = create_test_strategy();
3908 register_strategy(&mut strategy);
3909
3910 let orders = vec![
3911 make_initialized_market_order("O-20250208-LIST-CREATE-001"),
3912 make_initialized_market_order("O-20250208-LIST-CREATE-002"),
3913 ];
3914
3915 let client_order_id1 = orders[0].client_order_id();
3916 let client_order_id2 = orders[1].client_order_id();
3917 let cache_rc = strategy.core.cache_rc();
3918 let timeline = Rc::new(RefCell::new(Vec::new()));
3919 let event_messages = Rc::new(RefCell::new(Vec::new()));
3920
3921 let event_handler = {
3922 let event_messages = event_messages.clone();
3923 let timeline = timeline.clone();
3924 TypedHandler::from_with_id(
3925 "events.order.list_create_initialized",
3926 move |event: &OrderEventAny| {
3927 match event {
3928 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id1 => {
3929 let cache = cache_rc.borrow();
3930 let cached_order1 = cache.order(&client_order_id1).unwrap();
3931 let cached_order2 = cache.order(&client_order_id2).unwrap();
3932 let order_list_id = cached_order1.order_list_id().unwrap();
3933 assert_eq!(cached_order2.order_list_id(), Some(order_list_id));
3934 assert_eq!(e.order_list_id, Some(order_list_id));
3935 assert!(cache.order_list_exists(&order_list_id));
3936 let order_list = cache.order_list(&order_list_id).unwrap();
3937 assert_eq!(
3938 order_list.client_order_ids.as_slice(),
3939 &[client_order_id1, client_order_id2]
3940 );
3941 timeline.borrow_mut().push("init1");
3942 }
3943 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id2 => {
3944 let cache = cache_rc.borrow();
3945 let cached_order = cache.order(&client_order_id2).unwrap();
3946 assert_eq!(e.order_list_id, cached_order.order_list_id());
3947 timeline.borrow_mut().push("init2");
3948 }
3949 _ => panic!("unexpected order event {event:?}"),
3950 }
3951 event_messages.borrow_mut().push(event.clone());
3952 },
3953 )
3954 };
3955 let risk_handler = {
3956 let timeline = timeline.clone();
3957 TypedIntoHandler::from_with_id(
3958 "RiskEngine.queue_execute",
3959 move |command: TradingCommand| {
3960 let TradingCommand::SubmitOrderList(command) = command else {
3961 panic!("expected SubmitOrderList command");
3962 };
3963 assert!(
3964 command
3965 .order_inits
3966 .iter()
3967 .all(|init| init.order_list_id == Some(command.order_list.id))
3968 );
3969 timeline.borrow_mut().push("command");
3970 },
3971 )
3972 };
3973 msgbus::register_trading_command_endpoint(
3974 MessagingSwitchboard::risk_engine_queue_execute(),
3975 risk_handler,
3976 );
3977
3978 let topic = format!("events.order.{}", orders[0].strategy_id());
3979 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3980
3981 strategy
3982 .submit_order_list(orders, None, None, None)
3983 .unwrap();
3984
3985 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3986
3987 let cache = strategy.cache();
3988 let cached_order1 = cache.order(&client_order_id1).unwrap();
3989 let cached_order2 = cache.order(&client_order_id2).unwrap();
3990 let order_list_id = cached_order1.order_list_id().unwrap();
3991 assert_eq!(cached_order2.order_list_id(), Some(order_list_id));
3992
3993 let event_messages = event_messages.borrow();
3994 assert_eq!(event_messages.len(), 2);
3995 let OrderEventAny::Initialized(init1) = &event_messages[0] else {
3996 panic!("expected first OrderInitialized event");
3997 };
3998 let OrderEventAny::Initialized(init2) = &event_messages[1] else {
3999 panic!("expected second OrderInitialized event");
4000 };
4001 assert_eq!(init1.order_list_id, Some(order_list_id));
4002 assert_eq!(init2.order_list_id, Some(order_list_id));
4003 assert_eq!(timeline.borrow().as_slice(), &["init1", "init2", "command"]);
4004
4005 let order_list = cache.order_list(&order_list_id).unwrap();
4006 assert_eq!(
4007 order_list.client_order_ids.as_slice(),
4008 &[client_order_id1, client_order_id2]
4009 );
4010 }
4011
4012 #[rstest]
4013 fn test_submit_order_list_routes_optional_params_to_risk() {
4014 let mut strategy = create_test_strategy();
4015 register_strategy(&mut strategy);
4016
4017 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4018 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4019 msgbus::register_trading_command_endpoint(
4020 MessagingSwitchboard::risk_engine_queue_execute(),
4021 risk_handler,
4022 );
4023
4024 let no_params_orders = vec![
4025 make_initialized_market_order("O-20250208-LIST-001"),
4026 make_initialized_market_order("O-20250208-LIST-002"),
4027 ];
4028 strategy
4029 .submit_order_list(no_params_orders, None, None, None)
4030 .unwrap();
4031
4032 let mut params = Params::new();
4033 params.insert(
4034 "routing_hint".to_string(),
4035 Value::String("prefer_batch".to_string()),
4036 );
4037 let param_orders = vec![
4038 make_initialized_market_order("O-20250208-LIST-003"),
4039 make_initialized_market_order("O-20250208-LIST-004"),
4040 ];
4041 strategy
4042 .submit_order_list(param_orders, None, None, Some(params.clone()))
4043 .unwrap();
4044
4045 let risk_messages = risk_messages.get_messages();
4046 assert_eq!(risk_messages.len(), 2);
4047 let Some(TradingCommand::SubmitOrderList(no_params_command)) = risk_messages.first() else {
4048 panic!("expected SubmitOrderList command");
4049 };
4050 let Some(TradingCommand::SubmitOrderList(param_command)) = risk_messages.get(1) else {
4051 panic!("expected SubmitOrderList command");
4052 };
4053 assert!(no_params_command.params.is_none());
4054 assert_eq!(param_command.params.as_ref(), Some(¶ms));
4055 }
4056
4057 #[rstest]
4058 fn test_modify_order_routes_non_emulated_orders_to_risk() {
4059 let mut strategy = create_test_strategy();
4060 register_strategy(&mut strategy);
4061
4062 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4063 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4064 msgbus::register_trading_command_endpoint(
4065 MessagingSwitchboard::risk_engine_queue_execute(),
4066 risk_handler,
4067 );
4068
4069 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4070 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4071 msgbus::register_trading_command_endpoint(
4072 MessagingSwitchboard::exec_engine_queue_execute(),
4073 exec_handler,
4074 );
4075
4076 let order = OrderAny::Market(MarketOrder::new(
4077 TraderId::from("TRADER-001"),
4078 StrategyId::from("TEST-001"),
4079 InstrumentId::from("BTCUSDT.BINANCE"),
4080 ClientOrderId::from("O-20250208-0003"),
4081 OrderSide::Buy,
4082 Quantity::from(100_000),
4083 TimeInForce::Gtc,
4084 UUID4::new(),
4085 UnixNanos::default(),
4086 false,
4087 false,
4088 None,
4089 None,
4090 None,
4091 None,
4092 None,
4093 None,
4094 None,
4095 None,
4096 ));
4097 add_order_to_cache(&strategy, &order);
4098
4099 strategy
4100 .modify_order(
4101 order.client_order_id(),
4102 Some(Quantity::from(200_000)),
4103 None,
4104 None,
4105 None,
4106 None,
4107 )
4108 .unwrap();
4109
4110 let risk_messages = risk_messages.get_messages();
4111 let exec_messages = exec_messages.get_messages();
4112
4113 assert_eq!(risk_messages.len(), 1);
4114 assert!(matches!(
4115 risk_messages.first(),
4116 Some(TradingCommand::ModifyOrder(_))
4117 ));
4118 assert!(exec_messages.is_empty());
4119 }
4120
4121 #[rstest]
4122 fn test_modify_order_routes_active_local_algorithm_order_to_algorithm() {
4123 let mut strategy = create_test_strategy();
4124 register_strategy(&mut strategy);
4125
4126 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4127 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4128 msgbus::register_trading_command_endpoint(
4129 MessagingSwitchboard::risk_engine_queue_execute(),
4130 risk_handler,
4131 );
4132 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4133 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4134 msgbus::register_trading_command_endpoint(
4135 MessagingSwitchboard::exec_engine_queue_execute(),
4136 exec_handler,
4137 );
4138 let (algo_handler, algo_messages) =
4139 get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
4140 msgbus::register_any("TWAP.execute".into(), algo_handler);
4141
4142 let order = make_initialized_algorithm_order("O-20250208-ALGO-MODIFY-001");
4143 add_order_to_cache(&strategy, &order);
4144
4145 strategy
4146 .modify_order(
4147 order.client_order_id(),
4148 Some(Quantity::from(200_000)),
4149 None,
4150 None,
4151 None,
4152 None,
4153 )
4154 .unwrap();
4155
4156 let algo_messages = algo_messages.get_messages();
4157 assert_eq!(algo_messages.len(), 1);
4158 assert!(matches!(
4159 algo_messages.first(),
4160 Some(TradingCommand::ModifyOrder(command))
4161 if command.client_order_id == order.client_order_id()
4162 ));
4163 assert!(risk_messages.get_messages().is_empty());
4164 assert!(exec_messages.get_messages().is_empty());
4165 }
4166
4167 #[rstest]
4168 fn test_modify_order_routes_accepted_algorithm_order_to_risk() {
4169 let mut strategy = create_test_strategy();
4170 register_strategy(&mut strategy);
4171
4172 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4173 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4174 msgbus::register_trading_command_endpoint(
4175 MessagingSwitchboard::risk_engine_queue_execute(),
4176 risk_handler,
4177 );
4178 let (algo_handler, algo_messages) =
4179 get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
4180 msgbus::register_any("TWAP.execute".into(), algo_handler);
4181
4182 let mut order = make_initialized_algorithm_order("O-20250208-ALGO-MODIFY-002");
4183 let account_id = AccountId::from("ACC-001");
4184 order
4185 .apply(TestOrderEventStubs::submitted(&order, account_id))
4186 .unwrap();
4187 order
4188 .apply(TestOrderEventStubs::accepted(
4189 &order,
4190 account_id,
4191 VenueOrderId::from("O-20250208-ALGO-MODIFY-002"),
4192 ))
4193 .unwrap();
4194 add_order_to_cache(&strategy, &order);
4195
4196 strategy
4197 .modify_order(
4198 order.client_order_id(),
4199 Some(Quantity::from(200_000)),
4200 None,
4201 None,
4202 None,
4203 None,
4204 )
4205 .unwrap();
4206
4207 let risk_messages = risk_messages.get_messages();
4208 assert_eq!(risk_messages.len(), 1);
4209 assert!(matches!(
4210 risk_messages.first(),
4211 Some(TradingCommand::ModifyOrder(command))
4212 if command.client_order_id == order.client_order_id()
4213 ));
4214 assert!(algo_messages.get_messages().is_empty());
4215 assert_eq!(
4216 strategy
4217 .cache()
4218 .order(&order.client_order_id())
4219 .unwrap()
4220 .status(),
4221 OrderStatus::PendingUpdate
4222 );
4223 }
4224
4225 #[rstest]
4226 fn test_modify_order_routes_emulated_order_to_order_emulator() {
4227 let mut strategy = create_test_strategy();
4228 register_strategy(&mut strategy);
4229
4230 let (emulator_handler, emulator_messages): (
4231 _,
4232 TypedIntoMessageSavingHandler<TradingCommand>,
4233 ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4234 msgbus::register_trading_command_endpoint(
4235 MessagingSwitchboard::order_emulator_execute(),
4236 emulator_handler,
4237 );
4238 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4239 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4240 msgbus::register_trading_command_endpoint(
4241 MessagingSwitchboard::risk_engine_queue_execute(),
4242 risk_handler,
4243 );
4244 let mut order = OrderTestBuilder::new(OrderType::StopMarket)
4245 .trader_id(TraderId::from("TRADER-001"))
4246 .strategy_id(StrategyId::from("TEST-001"))
4247 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
4248 .client_order_id(ClientOrderId::from("O-20250208-EMULATED-MODIFY-001"))
4249 .side(OrderSide::Buy)
4250 .trigger_price(Price::from("51000.0"))
4251 .quantity(Quantity::from(100_000))
4252 .emulation_trigger(TriggerType::BidAsk)
4253 .build();
4254 order
4255 .apply(OrderEventAny::Emulated(
4256 OrderEmulatedSpec::builder()
4257 .trader_id(order.trader_id())
4258 .strategy_id(order.strategy_id())
4259 .instrument_id(order.instrument_id())
4260 .client_order_id(order.client_order_id())
4261 .build(),
4262 ))
4263 .unwrap();
4264 add_order_to_cache(&strategy, &order);
4265
4266 strategy
4267 .modify_order(
4268 order.client_order_id(),
4269 None,
4270 None,
4271 Some(Price::from("52000.0")),
4272 None,
4273 None,
4274 )
4275 .unwrap();
4276
4277 let emulator_messages = emulator_messages.get_messages();
4278 assert_eq!(emulator_messages.len(), 1);
4279 assert!(matches!(
4280 emulator_messages.first(),
4281 Some(TradingCommand::ModifyOrder(command))
4282 if command.client_order_id == order.client_order_id()
4283 ));
4284 assert!(risk_messages.get_messages().is_empty());
4285 }
4286
4287 #[rstest]
4288 fn test_modify_order_routes_initialized_trigger_order_to_emulator_reentrantly() {
4289 let strategy_id = StrategyId::from("REENTRANT-001");
4290 let modified_quantity = Quantity::from(200_000);
4291 let mut strategy = InitializedModifyStrategy {
4292 core: StrategyCore::new(StrategyConfig {
4293 strategy_id: Some(strategy_id),
4294 ..Default::default()
4295 }),
4296 modified_quantity,
4297 };
4298 let trader_id = TraderId::from("TRADER-001");
4299 let clock = Rc::new(RefCell::new(TestClock::new()));
4300 let cache = Rc::new(RefCell::new(Cache::default()));
4301 let portfolio = Rc::new(RefCell::new(Portfolio::new(
4302 clock.clone(),
4303 cache.clone(),
4304 None,
4305 )));
4306 strategy
4307 .core
4308 .register(trader_id, clock, cache, portfolio)
4309 .unwrap();
4310 strategy.initialize().unwrap();
4311 strategy.start().unwrap();
4312
4313 let actor_id = strategy.actor_id().inner();
4314 register_component_actor(strategy);
4315 let order_handler = TypedHandler::from(move |event: &OrderEventAny| {
4316 let mut strategy =
4317 try_get_actor_unchecked::<InitializedModifyStrategy>(&actor_id).unwrap();
4318 strategy.handle_order_event(event.clone());
4319 });
4320 let topic = format!("events.order.{strategy_id}");
4321 msgbus::subscribe_order_events(topic.clone().into(), order_handler.clone(), None);
4322
4323 let (emulator_handler, emulator_messages): (
4324 _,
4325 TypedIntoMessageSavingHandler<TradingCommand>,
4326 ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4327 msgbus::register_trading_command_endpoint(
4328 MessagingSwitchboard::order_emulator_execute(),
4329 emulator_handler,
4330 );
4331 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4332 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4333 msgbus::register_trading_command_endpoint(
4334 MessagingSwitchboard::risk_engine_queue_execute(),
4335 risk_handler,
4336 );
4337 let (algo_handler, algo_messages) =
4338 get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
4339 msgbus::register_any("TWAP.execute".into(), algo_handler);
4340
4341 let order = OrderTestBuilder::new(OrderType::StopMarket)
4342 .trader_id(trader_id)
4343 .strategy_id(strategy_id)
4344 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
4345 .client_order_id(ClientOrderId::from("O-REENTRANT-001"))
4346 .side(OrderSide::Buy)
4347 .trigger_price(Price::from("51000.0"))
4348 .quantity(Quantity::from(100_000))
4349 .emulation_trigger(TriggerType::BidAsk)
4350 .build();
4351 let client_order_id = order.client_order_id();
4352
4353 let mut strategy = try_get_actor_unchecked::<InitializedModifyStrategy>(&actor_id).unwrap();
4354 strategy.submit_order(order, None, None, None).unwrap();
4355 drop(strategy);
4356
4357 msgbus::unsubscribe_order_events(topic.into(), &order_handler);
4358 deregister_component(&strategy_id.inner());
4359 deregister_actor(&actor_id);
4360
4361 let emulator_messages = emulator_messages.get_messages();
4362 assert_eq!(emulator_messages.len(), 2);
4363 assert!(matches!(
4364 emulator_messages.first(),
4365 Some(TradingCommand::ModifyOrder(command))
4366 if command.client_order_id == client_order_id
4367 && command.quantity == Some(modified_quantity)
4368 ));
4369 assert!(
4370 risk_messages
4371 .get_messages()
4372 .iter()
4373 .all(|command| !matches!(command, TradingCommand::ModifyOrder(_)))
4374 );
4375 assert!(
4376 algo_messages
4377 .get_messages()
4378 .iter()
4379 .all(|command| !matches!(command, TradingCommand::ModifyOrder(_)))
4380 );
4381 assert!(matches!(
4382 emulator_messages.get(1),
4383 Some(TradingCommand::SubmitOrder(command))
4384 if command.client_order_id == client_order_id
4385 ));
4386 }
4387
4388 #[rstest]
4389 fn test_modify_order_prefers_emulator_over_algorithm_for_emulated_algorithm_order() {
4390 let mut strategy = create_test_strategy();
4391 register_strategy(&mut strategy);
4392
4393 let (emulator_handler, emulator_messages): (
4394 _,
4395 TypedIntoMessageSavingHandler<TradingCommand>,
4396 ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4397 msgbus::register_trading_command_endpoint(
4398 MessagingSwitchboard::order_emulator_execute(),
4399 emulator_handler,
4400 );
4401 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4402 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4403 msgbus::register_trading_command_endpoint(
4404 MessagingSwitchboard::risk_engine_queue_execute(),
4405 risk_handler,
4406 );
4407 let (algo_handler, algo_messages) =
4408 get_any_saving_handler::<TradingCommand>(Some(Ustr::from("TWAP.execute")));
4409 msgbus::register_any("TWAP.execute".into(), algo_handler);
4410
4411 let mut order = OrderTestBuilder::new(OrderType::StopMarket)
4412 .trader_id(TraderId::from("TRADER-001"))
4413 .strategy_id(StrategyId::from("TEST-001"))
4414 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
4415 .client_order_id(ClientOrderId::from("O-20250208-EMULATED-ALGO-MODIFY-001"))
4416 .side(OrderSide::Buy)
4417 .trigger_price(Price::from("51000.0"))
4418 .quantity(Quantity::from(100_000))
4419 .emulation_trigger(TriggerType::BidAsk)
4420 .exec_algorithm_id(ExecAlgorithmId::from("TWAP"))
4421 .exec_spawn_id(ClientOrderId::from("O-20250208-EMULATED-ALGO-MODIFY-001"))
4422 .build();
4423 order
4424 .apply(OrderEventAny::Emulated(
4425 OrderEmulatedSpec::builder()
4426 .trader_id(order.trader_id())
4427 .strategy_id(order.strategy_id())
4428 .instrument_id(order.instrument_id())
4429 .client_order_id(order.client_order_id())
4430 .build(),
4431 ))
4432 .unwrap();
4433
4434 assert!(order.is_emulated());
4437 assert!(order.is_active_local());
4438 assert!(order.exec_algorithm_id().is_some());
4439
4440 add_order_to_cache(&strategy, &order);
4441
4442 strategy
4443 .modify_order(
4444 order.client_order_id(),
4445 None,
4446 None,
4447 Some(Price::from("52000.0")),
4448 None,
4449 None,
4450 )
4451 .unwrap();
4452
4453 let emulator_messages = emulator_messages.get_messages();
4454 assert_eq!(emulator_messages.len(), 1);
4455 assert!(matches!(
4456 emulator_messages.first(),
4457 Some(TradingCommand::ModifyOrder(command))
4458 if command.client_order_id == order.client_order_id()
4459 ));
4460 assert!(algo_messages.get_messages().is_empty());
4461 assert!(risk_messages.get_messages().is_empty());
4462 }
4463
4464 #[rstest]
4465 fn test_modify_order_marks_order_pending_update_locally_before_send() {
4466 let mut strategy = create_test_strategy();
4467 register_strategy(&mut strategy);
4468
4469 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4470 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4471 msgbus::register_trading_command_endpoint(
4472 MessagingSwitchboard::risk_engine_queue_execute(),
4473 risk_handler,
4474 );
4475
4476 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
4477 get_typed_message_saving_handler(Some(Ustr::from("events.order.pending_update")));
4478 let order = make_accepted_limit_order("O-20250208-UPDATE-001");
4479 let topic = format!("events.order.{}", order.strategy_id());
4480 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4481 add_order_to_cache(&strategy, &order);
4482
4483 strategy
4484 .modify_order(
4485 order.client_order_id(),
4486 None,
4487 Some(Price::from("51000.0")),
4488 None,
4489 None,
4490 None,
4491 )
4492 .unwrap();
4493
4494 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4495
4496 let cache = strategy.cache();
4497 let cached_order = cache.order(&order.client_order_id()).unwrap();
4498 assert_eq!(cached_order.status(), OrderStatus::PendingUpdate);
4499
4500 let risk_messages = risk_messages.get_messages();
4501 assert_eq!(risk_messages.len(), 1);
4502 assert!(matches!(
4503 risk_messages.first(),
4504 Some(TradingCommand::ModifyOrder(_))
4505 ));
4506
4507 let event_messages = event_messages.get_messages();
4508 assert_eq!(event_messages.len(), 1);
4509 assert!(matches!(
4510 event_messages.first(),
4511 Some(OrderEventAny::PendingUpdate(_))
4512 ));
4513 }
4514
4515 #[rstest]
4516 fn test_modify_orders_marks_orders_pending_update_locally_before_send() {
4517 let mut strategy = create_test_strategy();
4518 register_strategy(&mut strategy);
4519
4520 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4521 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
4522 msgbus::register_trading_command_endpoint(
4523 MessagingSwitchboard::risk_engine_queue_execute(),
4524 risk_handler,
4525 );
4526
4527 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
4528 get_typed_message_saving_handler(Some(Ustr::from("events.order.batch_pending_update")));
4529 let order1 = make_accepted_limit_order("O-20250208-BATCH-UPDATE-001");
4530 let order2 = make_accepted_limit_order("O-20250208-BATCH-UPDATE-002");
4531 let topic = format!("events.order.{}", order1.strategy_id());
4532 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4533 add_order_to_cache(&strategy, &order1);
4534 add_order_to_cache(&strategy, &order2);
4535
4536 strategy
4537 .modify_orders(
4538 vec![
4539 (
4540 order1.client_order_id(),
4541 None,
4542 Some(Price::from("51000.0")),
4543 None,
4544 ),
4545 (
4546 order2.client_order_id(),
4547 Some(Quantity::from("2.0")),
4548 None,
4549 None,
4550 ),
4551 ],
4552 None,
4553 None,
4554 )
4555 .unwrap();
4556
4557 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4558
4559 let cache = strategy.cache();
4560 let cached_order1 = cache.order(&order1.client_order_id()).unwrap();
4561 let cached_order2 = cache.order(&order2.client_order_id()).unwrap();
4562 assert_eq!(cached_order1.status(), OrderStatus::PendingUpdate);
4563 assert_eq!(cached_order2.status(), OrderStatus::PendingUpdate);
4564
4565 let risk_messages = risk_messages.get_messages();
4566 assert_eq!(risk_messages.len(), 1);
4567 let Some(TradingCommand::ModifyOrders(command)) = risk_messages.first() else {
4568 panic!("expected BatchModifyOrders command");
4569 };
4570 assert_eq!(command.modifies.len(), 2);
4571 assert_eq!(
4572 command
4573 .modifies
4574 .iter()
4575 .map(|modify| modify.client_order_id)
4576 .collect::<Vec<_>>(),
4577 vec![order1.client_order_id(), order2.client_order_id()]
4578 );
4579
4580 let event_messages = event_messages.get_messages();
4581 assert_eq!(event_messages.len(), 2);
4582 assert!(
4583 event_messages
4584 .iter()
4585 .all(|event| matches!(event, OrderEventAny::PendingUpdate(_)))
4586 );
4587 }
4588
4589 #[rstest]
4590 fn test_cancel_order_marks_order_pending_cancel_locally_before_send() {
4591 let mut strategy = create_test_strategy();
4592 register_strategy(&mut strategy);
4593
4594 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4595 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4596 msgbus::register_trading_command_endpoint(
4597 MessagingSwitchboard::exec_engine_queue_execute(),
4598 exec_handler,
4599 );
4600
4601 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
4602 get_typed_message_saving_handler(Some(Ustr::from("events.order.pending_cancel")));
4603 let order = make_accepted_market_order("O-20250208-CANCEL-001");
4604 let topic = format!("events.order.{}", order.strategy_id());
4605 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4606 add_order_to_cache(&strategy, &order);
4607
4608 strategy
4609 .cancel_order(order.client_order_id(), None, None)
4610 .unwrap();
4611
4612 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4613
4614 let cache = strategy.cache();
4615 let cached_order = cache.order(&order.client_order_id()).unwrap();
4616 assert_eq!(cached_order.status(), OrderStatus::PendingCancel);
4617 let cache = strategy.core.cache_ref();
4618 assert!(cache.is_order_pending_cancel_local(&order.client_order_id()));
4619
4620 let exec_messages = exec_messages.get_messages();
4621 assert_eq!(exec_messages.len(), 1);
4622 assert!(matches!(
4623 exec_messages.first(),
4624 Some(TradingCommand::CancelOrder(_))
4625 ));
4626
4627 let event_messages = event_messages.get_messages();
4628 assert_eq!(event_messages.len(), 1);
4629 assert!(matches!(
4630 event_messages.first(),
4631 Some(OrderEventAny::PendingCancel(_))
4632 ));
4633 }
4634
4635 #[rstest]
4636 fn test_cancel_all_orders_strategy_only_sends_only_caller_strategy_cancels() {
4637 let mut strategy = create_test_strategy();
4638 register_strategy(&mut strategy);
4639
4640 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4641 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4642 msgbus::register_trading_command_endpoint(
4643 MessagingSwitchboard::exec_engine_queue_execute(),
4644 exec_handler,
4645 );
4646
4647 let order = make_accepted_market_order("O-20250208-CANCEL-ALL-001");
4648 let mut sibling_order = OrderTestBuilder::new(OrderType::Market)
4649 .trader_id(TraderId::from("TRADER-001"))
4650 .strategy_id(StrategyId::from("SIBLING-001"))
4651 .instrument_id(order.instrument_id())
4652 .client_order_id(ClientOrderId::from("O-20250208-CANCEL-ALL-002"))
4653 .side(OrderSide::Buy)
4654 .quantity(Quantity::from(100_000))
4655 .build();
4656 let account_id = AccountId::from("ACC-001");
4657 sibling_order
4658 .apply(TestOrderEventStubs::submitted(&sibling_order, account_id))
4659 .unwrap();
4660 sibling_order
4661 .apply(TestOrderEventStubs::accepted(
4662 &sibling_order,
4663 account_id,
4664 VenueOrderId::from("2"),
4665 ))
4666 .unwrap();
4667 add_order_to_cache(&strategy, &order);
4668 add_order_to_cache(&strategy, &sibling_order);
4669 strategy.core.cache_rc().borrow_mut().build_index();
4670
4671 strategy
4672 .cancel_all_orders(order.instrument_id(), None, None, true, None)
4673 .unwrap();
4674
4675 let messages = exec_messages.get_messages();
4676 let cache = strategy.cache();
4677 let cached_order = cache.order(&order.client_order_id()).unwrap();
4678 let cached_sibling = cache.order(&sibling_order.client_order_id()).unwrap();
4679 assert_eq!(messages.len(), 1);
4680 assert!(matches!(
4681 messages.first(),
4682 Some(TradingCommand::CancelOrder(command))
4683 if command.client_order_id == order.client_order_id()
4684 ));
4685 assert_eq!(cached_order.status(), OrderStatus::PendingCancel);
4686 assert_eq!(cached_sibling.status(), OrderStatus::Accepted);
4687 }
4688
4689 #[rstest]
4690 fn test_cancel_all_orders_strategy_only_deduplicates_emulated_inflight_order() {
4691 let mut strategy = create_test_strategy();
4692 register_strategy(&mut strategy);
4693
4694 let (emulator_handler, emulator_messages): (
4695 _,
4696 TypedIntoMessageSavingHandler<TradingCommand>,
4697 ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4698 msgbus::register_trading_command_endpoint(
4699 MessagingSwitchboard::order_emulator_execute(),
4700 emulator_handler,
4701 );
4702 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4703 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4704 msgbus::register_trading_command_endpoint(
4705 MessagingSwitchboard::exec_engine_queue_execute(),
4706 exec_handler,
4707 );
4708
4709 let client_order_id = ClientOrderId::from("O-20250208-CANCEL-ALL-EMULATED-001");
4710 let order = OrderTestBuilder::new(OrderType::StopMarket)
4711 .trader_id(TraderId::from("TRADER-001"))
4712 .strategy_id(StrategyId::from("TEST-001"))
4713 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
4714 .client_order_id(client_order_id)
4715 .side(OrderSide::Buy)
4716 .trigger_price(Price::from("51000.0"))
4717 .quantity(Quantity::from(100_000))
4718 .emulation_trigger(TriggerType::BidAsk)
4719 .exec_algorithm_id(ExecAlgorithmId::from("ALGO-001"))
4720 .exec_spawn_id(client_order_id)
4721 .build();
4722 add_order_to_cache(&strategy, &order);
4723 strategy.core.cache_rc().borrow_mut().build_index();
4724
4725 strategy
4726 .cancel_all_orders(order.instrument_id(), None, None, true, None)
4727 .unwrap();
4728
4729 let emulator_messages = emulator_messages.get_messages();
4730 assert_eq!(emulator_messages.len(), 1);
4731 assert!(matches!(
4732 emulator_messages.first(),
4733 Some(TradingCommand::CancelOrder(command))
4734 if command.client_order_id == order.client_order_id()
4735 ));
4736 assert!(exec_messages.get_messages().is_empty());
4737 }
4738
4739 #[rstest]
4740 fn test_cancel_all_orders_without_strategy_only_sends_cancel_all_command() {
4741 let mut strategy = create_test_strategy();
4742 register_strategy(&mut strategy);
4743
4744 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4745 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4746 msgbus::register_trading_command_endpoint(
4747 MessagingSwitchboard::exec_engine_queue_execute(),
4748 exec_handler,
4749 );
4750
4751 let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
4752
4753 strategy
4754 .cancel_all_orders(instrument_id, None, None, false, None)
4755 .unwrap();
4756
4757 let messages = exec_messages.get_messages();
4758 assert_eq!(messages.len(), 1);
4759 let Some(TradingCommand::CancelAllOrders(command)) = messages.first() else {
4760 panic!("Expected a CancelAllOrders command, was {messages:?}");
4761 };
4762 assert_eq!(command.strategy_id, StrategyId::from("TEST-001"));
4763 assert_eq!(command.instrument_id, instrument_id);
4764 assert_eq!(command.client_id, None);
4765 assert_eq!(command.order_side, None);
4766 assert_eq!(command.correlation_id, Some(command.command_id));
4767 assert_eq!(command.causation_id, None);
4768 }
4769
4770 #[rstest]
4771 fn test_cancel_all_orders_without_strategy_only_delegates_local_routing_with_params() {
4772 let mut strategy = create_test_strategy();
4773 register_strategy(&mut strategy);
4774
4775 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4776 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4777 msgbus::register_trading_command_endpoint(
4778 MessagingSwitchboard::exec_engine_queue_execute(),
4779 exec_handler,
4780 );
4781 let (emulator_handler, emulator_messages): (
4782 _,
4783 TypedIntoMessageSavingHandler<TradingCommand>,
4784 ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
4785 msgbus::register_trading_command_endpoint(
4786 MessagingSwitchboard::order_emulator_execute(),
4787 emulator_handler,
4788 );
4789 let exec_algorithm_id = ExecAlgorithmId::from("ALGO-001");
4790 let algorithm_endpoint = format!("{exec_algorithm_id}.execute");
4791 let (algorithm_handler, algorithm_messages) =
4792 get_any_saving_handler::<TradingCommand>(Some(Ustr::from("ALGO-001.execute")));
4793 msgbus::register_any(algorithm_endpoint.into(), algorithm_handler);
4794
4795 let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
4796 let sibling_strategy_id = StrategyId::from("SIBLING-001");
4797 let selected_client = ClientId::from("CLIENT-001");
4798 let account_id = AccountId::from("ACC-001");
4799
4800 let mut open_order = OrderTestBuilder::new(OrderType::Limit)
4801 .strategy_id(sibling_strategy_id)
4802 .instrument_id(instrument_id)
4803 .client_order_id(ClientOrderId::from("O-BROAD-OPEN-001"))
4804 .side(OrderSide::Buy)
4805 .price(Price::from("49000.0"))
4806 .quantity(Quantity::from(100_000))
4807 .build();
4808 open_order
4809 .apply(TestOrderEventStubs::submitted(&open_order, account_id))
4810 .unwrap();
4811 open_order
4812 .apply(TestOrderEventStubs::accepted(
4813 &open_order,
4814 account_id,
4815 VenueOrderId::from("V-BROAD-OPEN-001"),
4816 ))
4817 .unwrap();
4818 let emulated_order = OrderTestBuilder::new(OrderType::StopMarket)
4819 .strategy_id(sibling_strategy_id)
4820 .instrument_id(instrument_id)
4821 .client_order_id(ClientOrderId::from("O-BROAD-EMULATED-001"))
4822 .side(OrderSide::Buy)
4823 .trigger_price(Price::from("51000.0"))
4824 .quantity(Quantity::from(100_000))
4825 .emulation_trigger(TriggerType::BidAsk)
4826 .build();
4827 let algorithm_order_id = ClientOrderId::from("O-BROAD-ALGO-001");
4828 let algorithm_order = OrderTestBuilder::new(OrderType::Market)
4829 .strategy_id(sibling_strategy_id)
4830 .instrument_id(instrument_id)
4831 .client_order_id(algorithm_order_id)
4832 .side(OrderSide::Buy)
4833 .quantity(Quantity::from(100_000))
4834 .exec_algorithm_id(exec_algorithm_id)
4835 .exec_spawn_id(algorithm_order_id)
4836 .build();
4837
4838 {
4839 let cache_rc = strategy.core.cache_rc();
4840 let mut cache = cache_rc.borrow_mut();
4841 for order in [&open_order, &emulated_order, &algorithm_order] {
4842 cache
4843 .add_order(order.clone(), None, Some(selected_client), true)
4844 .unwrap();
4845 }
4846 cache.build_index();
4847 }
4848
4849 let mut params = Params::new();
4850 params.insert(
4851 "routing_hint".to_string(),
4852 Value::String("broad_cancel".to_string()),
4853 );
4854 strategy
4855 .cancel_all_orders(
4856 instrument_id,
4857 Some(OrderSide::Buy),
4858 Some(selected_client),
4859 false,
4860 Some(params.clone()),
4861 )
4862 .unwrap();
4863
4864 let exec_messages = exec_messages.get_messages();
4865 let Some(TradingCommand::CancelAllOrders(exec_command)) = exec_messages.first() else {
4866 panic!("Expected an execution CancelAllOrders command, was {exec_messages:?}");
4867 };
4868
4869 assert_eq!(exec_messages.len(), 1);
4870 assert!(emulator_messages.get_messages().is_empty());
4871 assert!(algorithm_messages.get_messages().is_empty());
4872 assert_eq!(exec_command.client_id, Some(selected_client));
4873 assert_eq!(exec_command.strategy_id, StrategyId::from("TEST-001"));
4874 assert_eq!(exec_command.instrument_id, instrument_id);
4875 assert_eq!(exec_command.order_side, Some(OrderSide::Buy));
4876 assert_eq!(exec_command.params.as_ref(), Some(¶ms));
4877 assert_eq!(exec_command.correlation_id, Some(exec_command.command_id));
4878 assert_eq!(exec_command.causation_id, None);
4879 }
4880
4881 #[rstest]
4882 fn test_cancel_all_orders_strategy_only_filters_by_client() {
4883 let mut strategy = create_test_strategy();
4884 register_strategy(&mut strategy);
4885
4886 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4887 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4888 msgbus::register_trading_command_endpoint(
4889 MessagingSwitchboard::exec_engine_queue_execute(),
4890 exec_handler,
4891 );
4892
4893 let selected_client = ClientId::from("CLIENT-001");
4894 let other_client = ClientId::from("CLIENT-002");
4895 let selected_order = make_accepted_market_order("O-20250208-CANCEL-CLIENT-001");
4896 let other_order = make_accepted_market_order("O-20250208-CANCEL-CLIENT-002");
4897 let cache_rc = strategy.core.cache_rc();
4898 {
4899 let mut cache = cache_rc.borrow_mut();
4900 cache
4901 .add_order(selected_order.clone(), None, Some(selected_client), true)
4902 .unwrap();
4903 cache
4904 .add_order(other_order.clone(), None, Some(other_client), true)
4905 .unwrap();
4906 cache.build_index();
4907 }
4908
4909 strategy
4910 .cancel_all_orders(
4911 selected_order.instrument_id(),
4912 None,
4913 Some(selected_client),
4914 true,
4915 None,
4916 )
4917 .unwrap();
4918
4919 let messages = exec_messages.get_messages();
4920 let cache = strategy.cache();
4921 let cached_selected = cache.order(&selected_order.client_order_id()).unwrap();
4922 let cached_other = cache.order(&other_order.client_order_id()).unwrap();
4923 assert_eq!(messages.len(), 1);
4924 assert!(matches!(
4925 messages.first(),
4926 Some(TradingCommand::CancelOrder(command))
4927 if command.client_id == Some(selected_client)
4928 && command.client_order_id == selected_order.client_order_id()
4929 ));
4930 assert_eq!(cached_selected.status(), OrderStatus::PendingCancel);
4931 assert_eq!(cached_other.status(), OrderStatus::Accepted);
4932 }
4933
4934 #[rstest]
4935 fn test_cancel_all_orders_strategy_only_filters_by_side_and_preserves_params() {
4936 let mut strategy = create_test_strategy();
4937 register_strategy(&mut strategy);
4938
4939 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
4940 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
4941 msgbus::register_trading_command_endpoint(
4942 MessagingSwitchboard::exec_engine_queue_execute(),
4943 exec_handler,
4944 );
4945
4946 let buy_order = make_accepted_market_order("O-20250208-CANCEL-SIDE-001");
4947 let mut sell_order = OrderTestBuilder::new(OrderType::Market)
4948 .trader_id(TraderId::from("TRADER-001"))
4949 .strategy_id(StrategyId::from("TEST-001"))
4950 .instrument_id(buy_order.instrument_id())
4951 .client_order_id(ClientOrderId::from("O-20250208-CANCEL-SIDE-002"))
4952 .side(OrderSide::Sell)
4953 .quantity(Quantity::from(100_000))
4954 .build();
4955 let account_id = AccountId::from("ACC-001");
4956 sell_order
4957 .apply(TestOrderEventStubs::submitted(&sell_order, account_id))
4958 .unwrap();
4959 sell_order
4960 .apply(TestOrderEventStubs::accepted(
4961 &sell_order,
4962 account_id,
4963 VenueOrderId::from("O-20250208-CANCEL-SIDE-002"),
4964 ))
4965 .unwrap();
4966 add_order_to_cache(&strategy, &buy_order);
4967 add_order_to_cache(&strategy, &sell_order);
4968 strategy.core.cache_rc().borrow_mut().build_index();
4969
4970 let mut params = Params::new();
4971 params.insert(
4972 "routing_hint".to_string(),
4973 Value::String("strategy_only".to_string()),
4974 );
4975 strategy
4976 .cancel_all_orders(
4977 buy_order.instrument_id(),
4978 Some(OrderSide::Buy),
4979 None,
4980 true,
4981 Some(params.clone()),
4982 )
4983 .unwrap();
4984
4985 let messages = exec_messages.get_messages();
4986 let Some(TradingCommand::CancelOrder(command)) = messages.first() else {
4987 panic!("Expected a CancelOrder command, was {messages:?}");
4988 };
4989 let cache = strategy.cache();
4990 let cached_buy = cache.order(&buy_order.client_order_id()).unwrap();
4991 let cached_sell = cache.order(&sell_order.client_order_id()).unwrap();
4992 assert_eq!(messages.len(), 1);
4993 assert_eq!(command.trader_id, TraderId::from("TRADER-001"));
4994 assert_eq!(command.client_id, None);
4995 assert_eq!(command.strategy_id, StrategyId::from("TEST-001"));
4996 assert_eq!(command.instrument_id, buy_order.instrument_id());
4997 assert_eq!(command.client_order_id, buy_order.client_order_id());
4998 assert_eq!(command.venue_order_id, buy_order.venue_order_id());
4999 assert_eq!(command.params.as_ref(), Some(¶ms));
5000 assert_eq!(cached_buy.status(), OrderStatus::PendingCancel);
5001 assert_eq!(cached_sell.status(), OrderStatus::Accepted);
5002 }
5003
5004 #[rstest]
5005 fn test_cancel_all_orders_strategy_only_continues_after_error_and_returns_first_error() {
5006 let mut strategy = create_test_strategy();
5007 register_strategy(&mut strategy);
5008
5009 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5010 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5011 msgbus::register_trading_command_endpoint(
5012 MessagingSwitchboard::exec_engine_queue_execute(),
5013 exec_handler,
5014 );
5015
5016 let mut failing_order = make_accepted_market_order("O-20250208-CANCEL-ERROR-001");
5017 let OrderAny::Market(order) = &mut failing_order else {
5018 panic!("Expected a MarketOrder");
5019 };
5020 order.account_id = None;
5021 let succeeding_order = make_accepted_market_order("O-20250208-CANCEL-ERROR-002");
5022 add_order_to_cache(&strategy, &failing_order);
5023 add_order_to_cache(&strategy, &succeeding_order);
5024 strategy.core.cache_rc().borrow_mut().build_index();
5025
5026 let error = strategy
5027 .cancel_all_orders(failing_order.instrument_id(), None, None, true, None)
5028 .unwrap_err()
5029 .to_string();
5030
5031 let messages = exec_messages.get_messages();
5032 let cache = strategy.cache();
5033 let cached_failing = cache.order(&failing_order.client_order_id()).unwrap();
5034 let cached_succeeding = cache.order(&succeeding_order.client_order_id()).unwrap();
5035 assert_eq!(
5036 error,
5037 "Cannot generate pending cancel event for O-20250208-CANCEL-ERROR-001: \
5038 account_id is not set"
5039 );
5040 assert_eq!(messages.len(), 1);
5041 assert!(matches!(
5042 messages.first(),
5043 Some(TradingCommand::CancelOrder(command))
5044 if command.client_order_id == succeeding_order.client_order_id()
5045 ));
5046 assert_eq!(cached_failing.status(), OrderStatus::Accepted);
5047 assert_eq!(cached_succeeding.status(), OrderStatus::PendingCancel);
5048 }
5049
5050 #[rstest]
5051 fn test_cancel_orders_marks_orders_pending_cancel_locally_before_send() {
5052 let mut strategy = create_test_strategy();
5053 register_strategy(&mut strategy);
5054
5055 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5056 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5057 msgbus::register_trading_command_endpoint(
5058 MessagingSwitchboard::exec_engine_queue_execute(),
5059 exec_handler,
5060 );
5061
5062 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
5063 get_typed_message_saving_handler(Some(Ustr::from("events.order.batch_pending_cancel")));
5064 let order1 = make_accepted_market_order("O-20250208-CANCEL-001");
5065 let order2 = make_accepted_market_order("O-20250208-CANCEL-002");
5066 let topic = format!("events.order.{}", order1.strategy_id());
5067 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
5068 add_order_to_cache(&strategy, &order1);
5069 add_order_to_cache(&strategy, &order2);
5070
5071 strategy
5072 .cancel_orders(
5073 vec![order1.client_order_id(), order2.client_order_id()],
5074 None,
5075 None,
5076 )
5077 .unwrap();
5078
5079 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
5080
5081 let cache = strategy.cache();
5082 let cached_order1 = cache.order(&order1.client_order_id()).unwrap();
5083 let cached_order2 = cache.order(&order2.client_order_id()).unwrap();
5084 assert_eq!(cached_order1.status(), OrderStatus::PendingCancel);
5085 assert_eq!(cached_order2.status(), OrderStatus::PendingCancel);
5086 let cache = strategy.core.cache_ref();
5087 assert!(cache.is_order_pending_cancel_local(&order1.client_order_id()));
5088 assert!(cache.is_order_pending_cancel_local(&order2.client_order_id()));
5089
5090 let exec_messages = exec_messages.get_messages();
5091 assert_eq!(exec_messages.len(), 1);
5092 let Some(TradingCommand::CancelOrders(command)) = exec_messages.first() else {
5093 panic!("expected BatchCancelOrders command");
5094 };
5095 assert_eq!(command.cancels.len(), 2);
5096
5097 let event_messages = event_messages.get_messages();
5098 assert_eq!(event_messages.len(), 2);
5099 assert!(
5100 event_messages
5101 .iter()
5102 .all(|event| matches!(event, OrderEventAny::PendingCancel(_)))
5103 );
5104 }
5105
5106 #[rstest]
5107 fn test_cancel_order_updates_own_book_status_before_send() {
5108 let mut strategy = create_test_strategy();
5109 register_strategy(&mut strategy);
5110
5111 let (exec_handler, _exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5112 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5113 msgbus::register_trading_command_endpoint(
5114 MessagingSwitchboard::exec_engine_queue_execute(),
5115 exec_handler,
5116 );
5117
5118 let order = make_accepted_limit_order("O-20250208-CANCEL-OWN-BOOK-001");
5119 add_order_to_cache_and_own_book(&strategy, &order);
5120
5121 strategy
5122 .cancel_order(order.client_order_id(), None, None)
5123 .unwrap();
5124
5125 let mut accepted = AHashSet::new();
5126 accepted.insert(OrderStatus::Accepted);
5127 let mut pending_cancel = AHashSet::new();
5128 pending_cancel.insert(OrderStatus::PendingCancel);
5129
5130 let cache = strategy.cache();
5131 let own_book = cache.own_order_book(&order.instrument_id()).unwrap();
5132 assert!(own_book.bids_as_map(Some(&accepted), None, None).is_empty());
5133 let pending_bids = own_book.bids_as_map(Some(&pending_cancel), None, None);
5134 assert_eq!(pending_bids.values().map(Vec::len).sum::<usize>(), 1);
5135 }
5136
5137 #[rstest]
5138 fn test_cancel_order_returns_error_when_not_in_cache() {
5139 let mut strategy = create_test_strategy();
5140 register_strategy(&mut strategy);
5141
5142 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5143 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5144 msgbus::register_trading_command_endpoint(
5145 MessagingSwitchboard::exec_engine_queue_execute(),
5146 exec_handler,
5147 );
5148
5149 let missing_id = ClientOrderId::from("O-MISSING");
5150 let err = strategy
5151 .cancel_order(missing_id, None, None)
5152 .expect_err("expected cancel_order to fail when order is not in cache");
5153
5154 assert_eq!(
5155 err.to_string(),
5156 format!("Cannot cancel order: {ORDER_NOT_FOUND}: {missing_id}")
5157 );
5158 assert!(exec_messages.get_messages().is_empty());
5159 }
5160
5161 #[rstest]
5162 fn test_modify_order_returns_error_when_not_in_cache() {
5163 let mut strategy = create_test_strategy();
5164 register_strategy(&mut strategy);
5165
5166 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5167 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
5168 msgbus::register_trading_command_endpoint(
5169 MessagingSwitchboard::risk_engine_queue_execute(),
5170 risk_handler,
5171 );
5172
5173 let missing_id = ClientOrderId::from("O-MISSING");
5174 let err = strategy
5175 .modify_order(missing_id, Some(Quantity::from(1)), None, None, None, None)
5176 .expect_err("expected modify_order to fail when order is not in cache");
5177
5178 assert_eq!(
5179 err.to_string(),
5180 format!("Cannot modify order: {ORDER_NOT_FOUND}: {missing_id}")
5181 );
5182 assert!(risk_messages.get_messages().is_empty());
5183 }
5184
5185 #[rstest]
5186 fn test_modify_orders_returns_error_when_any_id_missing() {
5187 let mut strategy = create_test_strategy();
5188 register_strategy(&mut strategy);
5189
5190 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5191 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
5192 msgbus::register_trading_command_endpoint(
5193 MessagingSwitchboard::risk_engine_queue_execute(),
5194 risk_handler,
5195 );
5196
5197 let order = make_accepted_limit_order("O-PRESENT");
5198 add_order_to_cache(&strategy, &order);
5199
5200 let missing_id = ClientOrderId::from("O-MISSING");
5201 let err = strategy
5202 .modify_orders(
5203 vec![
5204 (
5205 order.client_order_id(),
5206 None,
5207 Some(Price::from("51000.0")),
5208 None,
5209 ),
5210 (missing_id, Some(Quantity::from("2.0")), None, None),
5211 ],
5212 None,
5213 None,
5214 )
5215 .expect_err("expected modify_orders to fail when any id is missing");
5216
5217 assert_eq!(
5218 err.to_string(),
5219 format!("Cannot modify order: {ORDER_NOT_FOUND}: {missing_id}")
5220 );
5221 assert!(risk_messages.get_messages().is_empty());
5222 }
5223
5224 #[rstest]
5225 fn test_cancel_orders_returns_error_when_any_id_missing() {
5226 let mut strategy = create_test_strategy();
5227 register_strategy(&mut strategy);
5228
5229 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5230 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5231 msgbus::register_trading_command_endpoint(
5232 MessagingSwitchboard::exec_engine_queue_execute(),
5233 exec_handler,
5234 );
5235
5236 let order = make_accepted_limit_order("O-PRESENT");
5237 add_order_to_cache(&strategy, &order);
5238
5239 let missing_id = ClientOrderId::from("O-MISSING");
5240 let err = strategy
5241 .cancel_orders(vec![order.client_order_id(), missing_id], None, None)
5242 .expect_err("expected cancel_orders to fail when any id is missing");
5243
5244 assert_eq!(
5245 err.to_string(),
5246 format!("Cannot cancel order: {ORDER_NOT_FOUND}: {missing_id}")
5247 );
5248 assert!(exec_messages.get_messages().is_empty());
5249 }
5250
5251 #[rstest]
5254 fn test_has_gtd_expiry_timer_when_timer_not_set() {
5255 let mut strategy = create_test_strategy();
5256 let client_order_id = ClientOrderId::from("O-001");
5257
5258 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5259 }
5260
5261 #[rstest]
5262 fn test_has_gtd_expiry_timer_when_timer_set() {
5263 let mut strategy = create_test_strategy();
5264 let client_order_id = ClientOrderId::from("O-001");
5265
5266 strategy
5267 .core
5268 .gtd_timers
5269 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5270
5271 assert!(strategy.has_gtd_expiry_timer(&client_order_id));
5272 }
5273
5274 #[rstest]
5275 fn test_cancel_gtd_expiry_removes_timer() {
5276 let mut strategy = create_test_strategy();
5277 register_strategy(&mut strategy);
5278
5279 let client_order_id = ClientOrderId::from("O-001");
5280 strategy
5281 .core
5282 .gtd_timers
5283 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5284
5285 strategy.cancel_gtd_expiry(&client_order_id);
5286
5287 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5288 }
5289
5290 #[rstest]
5291 fn test_cancel_gtd_expiry_when_timer_not_set() {
5292 let mut strategy = create_test_strategy();
5293 register_strategy(&mut strategy);
5294
5295 let client_order_id = ClientOrderId::from("O-001");
5296
5297 strategy.cancel_gtd_expiry(&client_order_id);
5298
5299 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5300 }
5301
5302 #[rstest]
5303 #[case::matching_order(Ustr::from("GTD-EXPIRY:O-PRESENT"))]
5304 #[case::empty_order_id(Ustr::from("GTD-EXPIRY:"))]
5305 fn test_route_time_event_ignores_unregistered_gtd_timer(#[case] timer_name: Ustr) {
5306 let mut strategy = create_test_strategy();
5307 register_strategy(&mut strategy);
5308
5309 let order = make_accepted_limit_order("O-PRESENT");
5310 let client_order_id = order.client_order_id();
5311 add_order_to_cache(&strategy, &order);
5312 let event = TimeEvent::new(
5313 timer_name,
5314 UUID4::new(),
5315 UnixNanos::default(),
5316 UnixNanos::default(),
5317 );
5318
5319 route_time_event(&mut strategy, &event);
5320
5321 let cache = strategy.core.cache_ref();
5322 let cached_order = cache.order(&client_order_id).unwrap();
5323 assert_eq!(cached_order.status(), OrderStatus::Accepted);
5324 assert!(!strategy.core.gtd_timers.contains_key(&client_order_id));
5325 }
5326
5327 #[rstest]
5328 #[case::filled(make_filled)]
5329 #[case::canceled(make_canceled)]
5330 #[case::rejected(make_rejected)]
5331 #[case::expired(make_expired)]
5332 #[case::fill_voided(make_terminal_fill_voided)]
5333 fn test_handle_order_event_cancels_gtd_timer_for_terminal_event(
5334 #[case] make_event: fn(ClientOrderId) -> OrderEventAny,
5335 ) {
5336 let mut strategy = create_test_strategy();
5337 register_strategy(&mut strategy);
5338 start_strategy(&mut strategy);
5339
5340 let client_order_id = ClientOrderId::from("O-001");
5341 strategy
5342 .core
5343 .gtd_timers
5344 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5345
5346 strategy.handle_order_event(make_event(client_order_id));
5347
5348 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5349 }
5350
5351 #[rstest]
5352 #[case::partial_fill(make_filled)]
5353 #[case::non_reopened_fill_void(make_terminal_fill_voided)]
5354 fn test_handle_order_event_keeps_gtd_timer_when_cached_order_remains_open(
5355 #[case] make_event: fn(ClientOrderId) -> OrderEventAny,
5356 ) {
5357 let mut strategy = create_test_strategy();
5358 register_strategy(&mut strategy);
5359 start_strategy(&mut strategy);
5360
5361 let client_order_id = ClientOrderId::from("O-001");
5362 let order = make_accepted_limit_order(client_order_id.as_str());
5363 add_order_to_cache(&strategy, &order);
5364 strategy
5365 .core
5366 .gtd_timers
5367 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5368
5369 strategy.handle_order_event(make_event(client_order_id));
5370
5371 assert!(strategy.has_gtd_expiry_timer(&client_order_id));
5372 }
5373
5374 #[rstest]
5375 #[case::filled(make_filled)]
5376 #[case::canceled(make_canceled)]
5377 #[case::rejected(make_rejected)]
5378 #[case::expired(make_expired)]
5379 #[case::fill_voided(make_terminal_fill_voided)]
5380 fn test_handle_order_event_cancels_gtd_timer_when_stopped(
5381 #[case] make_event: fn(ClientOrderId) -> OrderEventAny,
5382 ) {
5383 let mut strategy = create_test_strategy();
5384 register_strategy(&mut strategy);
5385 start_strategy(&mut strategy);
5386
5387 let client_order_id = ClientOrderId::from("O-001");
5388 strategy
5389 .core
5390 .gtd_timers
5391 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5392
5393 stop_strategy(&mut strategy);
5394 assert_eq!(strategy.state(), ComponentState::Stopped);
5395
5396 strategy.handle_order_event(make_event(client_order_id));
5397
5398 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5399 }
5400
5401 #[rstest]
5402 fn test_handle_order_event_skips_gtd_cancel_for_non_terminal() {
5403 let mut strategy = create_test_strategy();
5404 register_strategy(&mut strategy);
5405 start_strategy(&mut strategy);
5406
5407 let client_order_id = ClientOrderId::from("O-001");
5408 strategy
5409 .core
5410 .gtd_timers
5411 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5412
5413 strategy.handle_order_event(make_accepted(client_order_id));
5414
5415 assert!(strategy.has_gtd_expiry_timer(&client_order_id));
5416 }
5417
5418 #[rstest]
5419 fn test_handle_reopened_fill_void_keeps_gtd_timer() {
5420 let mut strategy = create_test_strategy();
5421 register_strategy(&mut strategy);
5422 start_strategy(&mut strategy);
5423
5424 let client_order_id = ClientOrderId::from("O-001");
5425 strategy
5426 .core
5427 .gtd_timers
5428 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5429
5430 strategy.handle_order_event(make_fill_voided(client_order_id, true));
5431
5432 assert!(strategy.has_gtd_expiry_timer(&client_order_id));
5433 }
5434
5435 #[rstest]
5436 fn test_handle_cancel_rejected_restores_future_gtd_timer() {
5437 let mut strategy = create_gtd_managed_strategy();
5438 let clock = register_gtd_strategy(&mut strategy);
5439 start_strategy(&mut strategy);
5440 let expire_time = UnixNanos::from(clock.borrow().timestamp_ns().as_u64() + 1_000_000_000);
5441
5442 let order = make_accepted_gtd_limit_order("O-001", expire_time);
5443 let client_order_id = order.client_order_id();
5444 add_order_to_cache(&strategy, &order);
5445
5446 strategy.handle_order_event(make_cancel_rejected(client_order_id));
5447
5448 assert_eq!(
5449 strategy.core.gtd_timers.get(&client_order_id),
5450 Some(&Ustr::from("GTD-EXPIRY:O-001"))
5451 );
5452 }
5453
5454 #[rstest]
5455 fn test_cancel_rejected_restores_timer_after_cancel_order_removes_it() {
5456 let mut strategy = create_gtd_managed_strategy();
5457 let clock = register_gtd_strategy(&mut strategy);
5458 start_strategy(&mut strategy);
5459 let expire_time = UnixNanos::from(clock.borrow().timestamp_ns().as_u64() + 1_000_000_000);
5460 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5461 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5462 msgbus::register_trading_command_endpoint(
5463 MessagingSwitchboard::exec_engine_queue_execute(),
5464 exec_handler,
5465 );
5466
5467 let order = make_accepted_gtd_limit_order("O-001", expire_time);
5468 let client_order_id = order.client_order_id();
5469 add_order_to_cache(&strategy, &order);
5470 strategy
5471 .core
5472 .gtd_timers
5473 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5474
5475 strategy.cancel_order(client_order_id, None, None).unwrap();
5476 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5477
5478 let event = make_cancel_rejected(client_order_id);
5479 strategy
5480 .core
5481 .cache_rc()
5482 .borrow_mut()
5483 .update_order(&event)
5484 .unwrap();
5485 assert_eq!(
5486 strategy
5487 .core
5488 .cache_ref()
5489 .order(&client_order_id)
5490 .unwrap()
5491 .status(),
5492 OrderStatus::Accepted
5493 );
5494
5495 strategy.handle_order_event(event);
5496
5497 assert_eq!(
5498 strategy.core.gtd_timers.get(&client_order_id),
5499 Some(&Ustr::from("GTD-EXPIRY:O-001"))
5500 );
5501 assert_eq!(exec_messages.get_messages().len(), 1);
5502 }
5503
5504 #[rstest]
5505 fn test_cancel_rejected_while_submitted_restores_timer_through_acceptance() {
5506 let mut strategy = create_gtd_managed_strategy();
5507 let clock = register_gtd_strategy(&mut strategy);
5508 start_strategy(&mut strategy);
5509 let expire_time = UnixNanos::from(clock.borrow().timestamp_ns().as_u64() + 1_000_000_000);
5510 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5511 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5512 msgbus::register_trading_command_endpoint(
5513 MessagingSwitchboard::exec_engine_queue_execute(),
5514 exec_handler,
5515 );
5516
5517 let order = make_submitted_gtd_limit_order("O-001", expire_time);
5518 let client_order_id = order.client_order_id();
5519 add_order_to_cache(&strategy, &order);
5520 strategy.set_gtd_expiry(&order).unwrap();
5521
5522 strategy.cancel_order(client_order_id, None, None).unwrap();
5523 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5524
5525 let cancel_rejected = make_cancel_rejected(client_order_id);
5526 strategy
5527 .core
5528 .cache_rc()
5529 .borrow_mut()
5530 .update_order(&cancel_rejected)
5531 .unwrap();
5532 assert_eq!(
5533 strategy
5534 .core
5535 .cache_ref()
5536 .order(&client_order_id)
5537 .unwrap()
5538 .status(),
5539 OrderStatus::Submitted
5540 );
5541 strategy.handle_order_event(cancel_rejected);
5542 assert_eq!(
5543 strategy.core.gtd_timers.get(&client_order_id),
5544 Some(&Ustr::from("GTD-EXPIRY:O-001"))
5545 );
5546 assert_eq!(strategy.core.gtd_timers.len(), 1);
5547
5548 let accepted = make_accepted(client_order_id);
5549 strategy
5550 .core
5551 .cache_rc()
5552 .borrow_mut()
5553 .update_order(&accepted)
5554 .unwrap();
5555 strategy.handle_order_event(accepted);
5556
5557 assert_eq!(
5558 strategy.core.gtd_timers.get(&client_order_id),
5559 Some(&Ustr::from("GTD-EXPIRY:O-001"))
5560 );
5561 assert_eq!(strategy.core.gtd_timers.len(), 1);
5562 assert_eq!(exec_messages.get_messages().len(), 1);
5563 }
5564
5565 #[rstest]
5566 #[case::future_expiry(false)]
5567 #[case::elapsed_expiry(true)]
5568 fn test_cancel_rejected_after_acceptance_restores_gtd_expiry(#[case] expired: bool) {
5569 let mut strategy = create_gtd_managed_strategy();
5570 let clock = register_gtd_strategy(&mut strategy);
5571 start_strategy(&mut strategy);
5572 let expire_time = UnixNanos::from(1_000_000_000);
5573 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5574 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5575 msgbus::register_trading_command_endpoint(
5576 MessagingSwitchboard::exec_engine_queue_execute(),
5577 exec_handler,
5578 );
5579 let order = make_submitted_gtd_limit_order("O-001", expire_time);
5580 let client_order_id = order.client_order_id();
5581 add_order_to_cache(&strategy, &order);
5582 strategy.set_gtd_expiry(&order).unwrap();
5583
5584 strategy.cancel_order(client_order_id, None, None).unwrap();
5585 let accepted = make_accepted(client_order_id);
5586 strategy
5587 .core
5588 .cache_rc()
5589 .borrow_mut()
5590 .update_order(&accepted)
5591 .unwrap();
5592 strategy.handle_order_event(accepted);
5593 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5594 assert_eq!(exec_messages.get_messages().len(), 1);
5595
5596 if expired {
5597 clock.borrow_mut().set_time(expire_time);
5598 }
5599 let cancel_rejected = make_cancel_rejected(client_order_id);
5600 strategy
5601 .core
5602 .cache_rc()
5603 .borrow_mut()
5604 .update_order(&cancel_rejected)
5605 .unwrap();
5606 strategy.handle_order_event(cancel_rejected);
5607
5608 if !expired {
5609 assert_eq!(
5610 strategy.core.gtd_timers.get(&client_order_id),
5611 Some(&Ustr::from("GTD-EXPIRY:O-001"))
5612 );
5613 assert_eq!(
5614 strategy
5615 .core
5616 .cache_ref()
5617 .order(&client_order_id)
5618 .unwrap()
5619 .status(),
5620 OrderStatus::Accepted
5621 );
5622 assert_eq!(exec_messages.get_messages().len(), 1);
5623 let events = clock.borrow_mut().advance_time(expire_time, true);
5624 assert_eq!(events.len(), 1);
5625 route_time_event(&mut strategy, &events[0]);
5626 }
5627
5628 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5629 assert_eq!(
5630 strategy
5631 .core
5632 .cache_ref()
5633 .order(&client_order_id)
5634 .unwrap()
5635 .status(),
5636 OrderStatus::PendingCancel
5637 );
5638 let commands = exec_messages.get_messages();
5639 assert_eq!(commands.len(), 2);
5640 for command in commands {
5641 assert!(matches!(
5642 command,
5643 TradingCommand::CancelOrder(command) if command.client_order_id == client_order_id
5644 ));
5645 }
5646 }
5647
5648 #[rstest]
5649 fn test_cancel_rejected_while_submitted_retries_already_expired_gtd_order() {
5650 let mut strategy = create_gtd_managed_strategy();
5651 let clock = register_gtd_strategy(&mut strategy);
5652 start_strategy(&mut strategy);
5653 clock.borrow_mut().set_time(UnixNanos::from(1_000_000_000));
5654 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5655 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5656 msgbus::register_trading_command_endpoint(
5657 MessagingSwitchboard::exec_engine_queue_execute(),
5658 exec_handler,
5659 );
5660
5661 let order = make_submitted_gtd_limit_order("O-001", UnixNanos::from(1));
5662 let client_order_id = order.client_order_id();
5663 add_order_to_cache(&strategy, &order);
5664
5665 strategy.handle_order_event(make_cancel_rejected(client_order_id));
5666
5667 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5668 assert!(matches!(
5669 exec_messages.get_messages().as_slice(),
5670 [TradingCommand::CancelOrder(command)]
5671 if command.client_order_id == client_order_id
5672 ));
5673 assert_eq!(
5674 strategy
5675 .core
5676 .cache_ref()
5677 .order(&client_order_id)
5678 .unwrap()
5679 .status(),
5680 OrderStatus::PendingCancel
5681 );
5682 }
5683
5684 #[rstest]
5685 fn test_cancel_rejected_while_submitted_preserves_existing_gtd_timer() {
5686 let mut strategy = create_gtd_managed_strategy();
5687 let clock = register_gtd_strategy(&mut strategy);
5688 start_strategy(&mut strategy);
5689 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5690 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5691 msgbus::register_trading_command_endpoint(
5692 MessagingSwitchboard::exec_engine_queue_execute(),
5693 exec_handler,
5694 );
5695
5696 let order =
5697 make_submitted_gtd_limit_order("O-001", UnixNanos::from(2_000_000_000_000_000_000));
5698 let client_order_id = order.client_order_id();
5699 add_order_to_cache(&strategy, &order);
5700 let timer_name = Ustr::from("GTD-EXPIRY:O-001");
5701 strategy.core.gtd_timers.insert(client_order_id, timer_name);
5702
5703 strategy.handle_order_event(make_cancel_rejected(client_order_id));
5704
5705 assert_eq!(
5706 strategy.core.gtd_timers.get(&client_order_id),
5707 Some(&timer_name)
5708 );
5709 assert_eq!(strategy.core.gtd_timers.len(), 1);
5710 assert!(!clock.borrow().timer_names().contains(&timer_name.as_str()));
5713 assert!(exec_messages.get_messages().is_empty());
5714 }
5715
5716 #[rstest]
5717 fn test_handle_cancel_rejected_retries_already_expired_gtd_order() {
5718 let mut strategy = create_gtd_managed_strategy();
5719 let clock = register_gtd_strategy(&mut strategy);
5720 start_strategy(&mut strategy);
5721 clock.borrow_mut().set_time(UnixNanos::from(1_000_000_000));
5722 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5723 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5724 msgbus::register_trading_command_endpoint(
5725 MessagingSwitchboard::exec_engine_queue_execute(),
5726 exec_handler,
5727 );
5728
5729 let order = make_accepted_gtd_limit_order("O-001", UnixNanos::from(1));
5730 let client_order_id = order.client_order_id();
5731 add_order_to_cache(&strategy, &order);
5732
5733 strategy.handle_order_event(make_cancel_rejected(client_order_id));
5734
5735 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5736 assert!(matches!(
5737 exec_messages.get_messages().as_slice(),
5738 [TradingCommand::CancelOrder(command)]
5739 if command.client_order_id == client_order_id
5740 ));
5741 assert_eq!(
5742 strategy
5743 .core
5744 .cache_ref()
5745 .order(&client_order_id)
5746 .unwrap()
5747 .status(),
5748 OrderStatus::PendingCancel
5749 );
5750 }
5751
5752 #[rstest]
5753 fn test_handle_cancel_rejected_preserves_existing_gtd_timer() {
5754 let mut strategy = create_gtd_managed_strategy();
5755 let _clock = register_gtd_strategy(&mut strategy);
5756 start_strategy(&mut strategy);
5757 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5758 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5759 msgbus::register_trading_command_endpoint(
5760 MessagingSwitchboard::exec_engine_queue_execute(),
5761 exec_handler,
5762 );
5763
5764 let order =
5765 make_accepted_gtd_limit_order("O-001", UnixNanos::from(2_000_000_000_000_000_000));
5766 let client_order_id = order.client_order_id();
5767 add_order_to_cache(&strategy, &order);
5768 strategy
5769 .core
5770 .gtd_timers
5771 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5772
5773 strategy.handle_order_event(make_cancel_rejected(client_order_id));
5774
5775 assert_eq!(
5776 strategy.core.gtd_timers.get(&client_order_id),
5777 Some(&Ustr::from("GTD-EXPIRY:O-001"))
5778 );
5779 assert!(exec_messages.get_messages().is_empty());
5780 }
5781
5782 #[rstest]
5783 fn test_handle_cancel_rejected_preserves_existing_timer_for_expired_order() {
5784 let mut strategy = create_gtd_managed_strategy();
5785 let clock = register_gtd_strategy(&mut strategy);
5786 start_strategy(&mut strategy);
5787 clock.borrow_mut().set_time(UnixNanos::from(1_000_000_000));
5788 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5789 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5790 msgbus::register_trading_command_endpoint(
5791 MessagingSwitchboard::exec_engine_queue_execute(),
5792 exec_handler,
5793 );
5794
5795 let order = make_accepted_gtd_limit_order("O-001", UnixNanos::from(1));
5796 let client_order_id = order.client_order_id();
5797 add_order_to_cache(&strategy, &order);
5798 strategy
5799 .core
5800 .gtd_timers
5801 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
5802
5803 strategy.handle_order_event(make_cancel_rejected(client_order_id));
5804
5805 assert_eq!(
5806 strategy.core.gtd_timers.get(&client_order_id),
5807 Some(&Ustr::from("GTD-EXPIRY:O-001"))
5808 );
5809 assert!(exec_messages.get_messages().is_empty());
5810 }
5811
5812 #[rstest]
5813 fn test_restored_gtd_timer_event_cancels_order() {
5814 let mut strategy = create_gtd_managed_strategy();
5815 let clock = register_gtd_strategy(&mut strategy);
5816 start_strategy(&mut strategy);
5817 let expire_time = UnixNanos::from(clock.borrow().timestamp_ns().as_u64() + 1_000_000_000);
5818 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5819 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5820 msgbus::register_trading_command_endpoint(
5821 MessagingSwitchboard::exec_engine_queue_execute(),
5822 exec_handler,
5823 );
5824
5825 let order = make_accepted_gtd_limit_order("O-001", expire_time);
5826 let client_order_id = order.client_order_id();
5827 add_order_to_cache(&strategy, &order);
5828
5829 strategy.handle_order_event(make_cancel_rejected(client_order_id));
5830 assert!(strategy.has_gtd_expiry_timer(&client_order_id));
5831
5832 let events = clock.borrow_mut().advance_time(expire_time, true);
5833 for event in &events {
5834 route_time_event(&mut strategy, event);
5835 }
5836
5837 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5838 assert!(matches!(
5839 exec_messages.get_messages().as_slice(),
5840 [TradingCommand::CancelOrder(command)]
5841 if command.client_order_id == client_order_id
5842 ));
5843 }
5844
5845 enum IneligibleRestoreCase {
5846 MissingOrder,
5847 ClosedOrder,
5848 NonGtdOrder,
5849 ManagementDisabled,
5850 Stopped,
5851 Emulated,
5852 }
5853
5854 fn assert_cancel_rejected_does_not_restore(case: &IneligibleRestoreCase) {
5855 let mut strategy = match case {
5856 IneligibleRestoreCase::ManagementDisabled => create_test_strategy(),
5857 _ => create_gtd_managed_strategy(),
5858 };
5859 let _clock = register_gtd_strategy(&mut strategy);
5860 start_strategy(&mut strategy);
5861 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5862 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
5863 msgbus::register_trading_command_endpoint(
5864 MessagingSwitchboard::exec_engine_queue_execute(),
5865 exec_handler,
5866 );
5867
5868 let client_order_id = ClientOrderId::from("O-001");
5869 let order = match case {
5870 IneligibleRestoreCase::MissingOrder => None,
5871 IneligibleRestoreCase::NonGtdOrder => {
5872 Some(make_accepted_limit_order(client_order_id.as_str()))
5873 }
5874 IneligibleRestoreCase::ClosedOrder => {
5875 let mut order = make_accepted_gtd_limit_order(
5876 client_order_id.as_str(),
5877 UnixNanos::from(2_000_000_000_000_000_000),
5878 );
5879 order.apply(make_canceled(client_order_id)).unwrap();
5880 Some(order)
5881 }
5882 IneligibleRestoreCase::Emulated => {
5883 let mut order = make_submitted_gtd_limit_order(
5884 client_order_id.as_str(),
5885 UnixNanos::from(2_000_000_000_000_000_000),
5886 );
5887 order.set_emulation_trigger(Some(TriggerType::BidAsk));
5888 Some(order)
5889 }
5890 IneligibleRestoreCase::ManagementDisabled | IneligibleRestoreCase::Stopped => {
5891 Some(make_accepted_gtd_limit_order(
5892 client_order_id.as_str(),
5893 UnixNanos::from(2_000_000_000_000_000_000),
5894 ))
5895 }
5896 };
5897
5898 if let Some(order) = order {
5899 add_order_to_cache(&strategy, &order);
5900 }
5901
5902 if matches!(case, IneligibleRestoreCase::Stopped) {
5903 stop_strategy(&mut strategy);
5904 }
5905
5906 strategy.handle_order_event(make_cancel_rejected(client_order_id));
5907
5908 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
5909 assert!(exec_messages.get_messages().is_empty());
5910 }
5911
5912 #[rstest]
5913 fn test_handle_cancel_rejected_does_not_restore_missing_order() {
5914 assert_cancel_rejected_does_not_restore(&IneligibleRestoreCase::MissingOrder);
5915 }
5916
5917 #[rstest]
5918 fn test_handle_cancel_rejected_does_not_restore_closed_order() {
5919 assert_cancel_rejected_does_not_restore(&IneligibleRestoreCase::ClosedOrder);
5920 }
5921
5922 #[rstest]
5923 fn test_handle_cancel_rejected_does_not_restore_non_gtd_order() {
5924 assert_cancel_rejected_does_not_restore(&IneligibleRestoreCase::NonGtdOrder);
5925 }
5926
5927 #[rstest]
5928 fn test_handle_cancel_rejected_does_not_restore_when_disabled() {
5929 assert_cancel_rejected_does_not_restore(&IneligibleRestoreCase::ManagementDisabled);
5930 }
5931
5932 #[rstest]
5933 fn test_handle_cancel_rejected_does_not_restore_when_stopped() {
5934 assert_cancel_rejected_does_not_restore(&IneligibleRestoreCase::Stopped);
5935 }
5936
5937 #[rstest]
5938 fn test_handle_cancel_rejected_does_not_restore_emulated_order() {
5939 assert_cancel_rejected_does_not_restore(&IneligibleRestoreCase::Emulated);
5940 }
5941
5942 #[rstest]
5943 fn test_handle_order_event_skips_dispatch_when_stopped() {
5944 let mut strategy = create_test_strategy();
5945 register_strategy(&mut strategy);
5946 start_strategy(&mut strategy);
5947 stop_strategy(&mut strategy);
5948 assert_eq!(strategy.state(), ComponentState::Stopped);
5949
5950 strategy.handle_order_event(make_rejected(ClientOrderId::from("O-001")));
5951
5952 assert!(!strategy.on_order_event_called);
5953 assert!(!strategy.on_order_rejected_called);
5954 }
5955
5956 #[rstest]
5957 fn test_on_start_calls_reactivate_gtd_timers_when_enabled() {
5958 let config = StrategyConfig {
5959 strategy_id: Some(StrategyId::from("TEST-001")),
5960 order_id_tag: Some("001".to_string()),
5961 manage_gtd_expiry: true,
5962 ..Default::default()
5963 };
5964 let mut strategy = TestStrategy::new(config);
5965 register_strategy(&mut strategy);
5966
5967 let result = Strategy::on_start(&mut strategy);
5968 assert!(result.is_ok());
5969 }
5970
5971 #[rstest]
5972 fn test_on_start_does_not_panic_when_gtd_disabled() {
5973 let config = StrategyConfig {
5974 strategy_id: Some(StrategyId::from("TEST-001")),
5975 order_id_tag: Some("001".to_string()),
5976 manage_gtd_expiry: false,
5977 ..Default::default()
5978 };
5979 let mut strategy = TestStrategy::new(config);
5980 register_strategy(&mut strategy);
5981
5982 let result = Strategy::on_start(&mut strategy);
5983 assert!(result.is_ok());
5984 }
5985
5986 #[rstest]
5987 fn test_on_start_errors_when_strategy_id_is_not_set() {
5988 let mut strategy = TestStrategy::new(StrategyConfig::default());
5989
5990 let err = Strategy::on_start(&mut strategy).unwrap_err().to_string();
5991
5992 assert_eq!(err, "Strategy not registered: strategy_id is not set");
5993 }
5994
5995 #[rstest]
5998 fn test_query_account_when_registered() {
5999 let mut strategy = create_test_strategy();
6000 register_strategy(&mut strategy);
6001
6002 let account_id = AccountId::from("ACC-001");
6003
6004 let result = strategy.query_account(account_id, None, None);
6005
6006 assert!(result.is_ok());
6007 }
6008
6009 #[rstest]
6010 fn test_query_account_with_client_id() {
6011 let mut strategy = create_test_strategy();
6012 register_strategy(&mut strategy);
6013
6014 let account_id = AccountId::from("ACC-001");
6015 let client_id = ClientId::from("BINANCE");
6016
6017 let result = strategy.query_account(account_id, Some(client_id), None);
6018
6019 assert!(result.is_ok());
6020 }
6021
6022 #[rstest]
6023 fn test_query_order_when_registered() {
6024 let mut strategy = create_test_strategy();
6025 register_strategy(&mut strategy);
6026
6027 let order = OrderAny::Market(MarketOrder::test_default());
6028
6029 let result = strategy.query_order(&order, None, None);
6030
6031 assert!(result.is_ok());
6032 }
6033
6034 #[rstest]
6035 fn test_query_order_with_client_id() {
6036 let mut strategy = create_test_strategy();
6037 register_strategy(&mut strategy);
6038
6039 let order = OrderAny::Market(MarketOrder::test_default());
6040 let client_id = ClientId::from("BINANCE");
6041
6042 let result = strategy.query_order(&order, Some(client_id), None);
6043
6044 assert!(result.is_ok());
6045 }
6046
6047 #[rstest]
6048 fn test_is_exiting_returns_false_by_default() {
6049 let strategy = create_test_strategy();
6050 assert!(!strategy.is_exiting());
6051 }
6052
6053 #[rstest]
6054 fn test_is_exiting_returns_true_when_set_manually() {
6055 let mut strategy = create_test_strategy();
6056 register_strategy(&mut strategy);
6057
6058 strategy.core.is_exiting = true;
6060
6061 assert!(strategy.is_exiting());
6062 }
6063
6064 #[rstest]
6065 fn test_market_exit_sets_is_exiting_flag() {
6066 let mut strategy = create_test_strategy();
6068 register_strategy(&mut strategy);
6069
6070 assert!(!strategy.core.is_exiting);
6071
6072 strategy.core.is_exiting = true;
6074 strategy.core.market_exit_attempts = 0;
6075
6076 assert!(strategy.core.is_exiting);
6077 assert_eq!(strategy.core.market_exit_attempts, 0);
6078 }
6079
6080 #[rstest]
6081 fn test_market_exit_uses_config_time_in_force_and_reduce_only() {
6082 let config = StrategyConfig {
6083 strategy_id: Some(StrategyId::from("TEST-001")),
6084 order_id_tag: Some("001".to_string()),
6085 market_exit_time_in_force: TimeInForce::Ioc,
6086 market_exit_reduce_only: false,
6087 ..Default::default()
6088 };
6089 let strategy = TestStrategy::new(config);
6090
6091 assert_eq!(
6092 strategy.core.config.market_exit_time_in_force,
6093 TimeInForce::Ioc
6094 );
6095 assert!(!strategy.core.config.market_exit_reduce_only);
6096 }
6097
6098 #[rstest]
6099 fn test_market_exit_resets_attempt_counter() {
6100 let mut strategy = create_test_strategy();
6101 register_strategy(&mut strategy);
6102
6103 strategy.core.market_exit_attempts = 50;
6105
6106 strategy.core.reset_market_exit_state();
6108
6109 assert_eq!(strategy.core.market_exit_attempts, 0);
6110 }
6111
6112 #[rstest]
6113 fn test_market_exit_second_call_returns_early_when_exiting() {
6114 let mut strategy = create_test_strategy();
6115 register_strategy(&mut strategy);
6116
6117 strategy.core.is_exiting = true;
6119
6120 let result = strategy.market_exit();
6122 assert!(result.is_ok());
6123 assert!(strategy.core.is_exiting);
6124 }
6125
6126 #[rstest]
6127 fn test_finalize_market_exit_resets_state() {
6128 let mut strategy = create_test_strategy();
6129 register_strategy(&mut strategy);
6130
6131 strategy.core.is_exiting = true;
6133 strategy.core.pending_stop = true;
6134 strategy.core.market_exit_attempts = 50;
6135
6136 strategy.finalize_market_exit();
6137
6138 assert!(!strategy.core.is_exiting);
6139 assert!(!strategy.core.pending_stop);
6140 assert_eq!(strategy.core.market_exit_attempts, 0);
6141 }
6142
6143 #[rstest]
6144 fn test_market_exit_config_defaults() {
6145 let config = StrategyConfig::default();
6146
6147 assert!(!config.manage_stop);
6148 assert_eq!(config.market_exit_interval_ms, 100);
6149 assert_eq!(config.market_exit_max_attempts, 100);
6150 }
6151
6152 #[rstest]
6153 fn test_market_exit_with_custom_config() {
6154 let config = StrategyConfig {
6155 strategy_id: Some(StrategyId::from("TEST-001")),
6156 manage_stop: true,
6157 market_exit_interval_ms: 50,
6158 market_exit_max_attempts: 200,
6159 ..Default::default()
6160 };
6161 let strategy = TestStrategy::new(config);
6162
6163 assert!(strategy.core.config.manage_stop);
6164 assert_eq!(strategy.core.config.market_exit_interval_ms, 50);
6165 assert_eq!(strategy.core.config.market_exit_max_attempts, 200);
6166 }
6167
6168 #[derive(Debug)]
6169 struct MarketExitHookTrackingStrategy {
6170 core: StrategyCore,
6171 on_market_exit_called: bool,
6172 post_market_exit_called: bool,
6173 }
6174
6175 impl MarketExitHookTrackingStrategy {
6176 fn new(config: StrategyConfig) -> Self {
6177 Self {
6178 core: StrategyCore::new(config),
6179 on_market_exit_called: false,
6180 post_market_exit_called: false,
6181 }
6182 }
6183 }
6184
6185 impl DataActor for MarketExitHookTrackingStrategy {}
6186
6187 nautilus_strategy!(MarketExitHookTrackingStrategy, {
6188 fn on_market_exit(&mut self) {
6189 self.on_market_exit_called = true;
6190 }
6191
6192 fn post_market_exit(&mut self) {
6193 self.post_market_exit_called = true;
6194 }
6195 });
6196
6197 #[rstest]
6198 fn test_market_exit_calls_on_market_exit_hook() {
6199 let config = StrategyConfig {
6200 strategy_id: Some(StrategyId::from("TEST-001")),
6201 order_id_tag: Some("001".to_string()),
6202 ..Default::default()
6203 };
6204 let mut strategy = MarketExitHookTrackingStrategy::new(config);
6205
6206 let trader_id = TraderId::from("TRADER-001");
6207 let clock = Rc::new(RefCell::new(TestClock::new()));
6208 let cache = Rc::new(RefCell::new(Cache::default()));
6209 let portfolio = Rc::new(RefCell::new(Portfolio::new(
6210 clock.clone(),
6211 cache.clone(),
6212 None,
6213 )));
6214 strategy
6215 .core
6216 .register(trader_id, clock, cache, portfolio)
6217 .unwrap();
6218 strategy.initialize().unwrap();
6219 strategy.start().unwrap();
6220
6221 let _ = strategy.market_exit();
6222
6223 assert!(strategy.on_market_exit_called);
6224 }
6225
6226 #[rstest]
6227 fn test_finalize_market_exit_calls_post_market_exit_hook() {
6228 let config = StrategyConfig {
6229 strategy_id: Some(StrategyId::from("TEST-001")),
6230 order_id_tag: Some("001".to_string()),
6231 ..Default::default()
6232 };
6233 let mut strategy = MarketExitHookTrackingStrategy::new(config);
6234
6235 let trader_id = TraderId::from("TRADER-001");
6236 let clock = Rc::new(RefCell::new(TestClock::new()));
6237 let cache = Rc::new(RefCell::new(Cache::default()));
6238 let portfolio = Rc::new(RefCell::new(Portfolio::new(
6239 clock.clone(),
6240 cache.clone(),
6241 None,
6242 )));
6243 strategy
6244 .core
6245 .register(trader_id, clock, cache, portfolio)
6246 .unwrap();
6247
6248 strategy.core.is_exiting = true;
6249 strategy.finalize_market_exit();
6250
6251 assert!(strategy.post_market_exit_called);
6252 }
6253
6254 #[derive(Debug)]
6255 struct FailingPostExitStrategy {
6256 core: StrategyCore,
6257 }
6258
6259 impl FailingPostExitStrategy {
6260 fn new(config: StrategyConfig) -> Self {
6261 Self {
6262 core: StrategyCore::new(config),
6263 }
6264 }
6265 }
6266
6267 impl DataActor for FailingPostExitStrategy {}
6268
6269 nautilus_strategy!(FailingPostExitStrategy, {
6270 fn post_market_exit(&mut self) {
6271 panic!("Simulated error in post_market_exit");
6272 }
6273 });
6274
6275 #[rstest]
6276 fn test_finalize_market_exit_handles_hook_panic() {
6277 let config = StrategyConfig {
6278 strategy_id: Some(StrategyId::from("TEST-001")),
6279 order_id_tag: Some("001".to_string()),
6280 ..Default::default()
6281 };
6282 let mut strategy = FailingPostExitStrategy::new(config);
6283
6284 let trader_id = TraderId::from("TRADER-001");
6285 let clock = Rc::new(RefCell::new(TestClock::new()));
6286 let cache = Rc::new(RefCell::new(Cache::default()));
6287 let portfolio = Rc::new(RefCell::new(Portfolio::new(
6288 clock.clone(),
6289 cache.clone(),
6290 None,
6291 )));
6292 strategy
6293 .core
6294 .register(trader_id, clock, cache, portfolio)
6295 .unwrap();
6296
6297 strategy.core.is_exiting = true;
6298 strategy.core.pending_stop = true;
6299
6300 strategy.finalize_market_exit();
6302
6303 assert!(!strategy.core.is_exiting);
6305 assert!(!strategy.core.pending_stop);
6306 }
6307
6308 #[rstest]
6309 fn test_check_market_exit_increments_attempts_before_finalizing() {
6310 let mut strategy = create_test_strategy();
6311 register_strategy(&mut strategy);
6312
6313 strategy.core.is_exiting = true;
6314 assert_eq!(strategy.core.market_exit_attempts, 0);
6315
6316 let event = TimeEvent::new(
6317 Ustr::from("MARKET_EXIT_CHECK:TEST-001"),
6318 UUID4::new(),
6319 UnixNanos::default(),
6320 UnixNanos::default(),
6321 );
6322 strategy.check_market_exit(event);
6323
6324 assert!(!strategy.core.is_exiting);
6328 assert_eq!(strategy.core.market_exit_attempts, 0);
6329 }
6330
6331 #[rstest]
6332 fn test_route_time_event_is_idempotent_when_callback_forwards() {
6333 let mut strategy = TestStrategy::new(StrategyConfig {
6334 strategy_id: Some(StrategyId::from("TEST-001")),
6335 ..Default::default()
6336 });
6337 register_strategy(&mut strategy);
6338
6339 let order = make_accepted_limit_order("O-PRESENT");
6340 add_order_to_cache(&strategy, &order);
6341 strategy.core.cache_rc().borrow_mut().build_index();
6342 strategy.core.is_exiting = true;
6343 let event = TimeEvent::new(
6344 strategy.core.market_exit_timer_name,
6345 UUID4::new(),
6346 UnixNanos::default(),
6347 UnixNanos::default(),
6348 );
6349
6350 route_time_event(&mut strategy, &event);
6351 Strategy::on_time_event(&mut strategy, &event).unwrap();
6352
6353 assert!(strategy.core.is_exiting);
6354 assert_eq!(strategy.core.market_exit_attempts, 1);
6355 assert_eq!(
6356 strategy.core.managed_time_event_last_id,
6357 Some(event.event_id)
6358 );
6359 }
6360
6361 #[rstest]
6362 fn test_route_time_event_deduplicates_overridden_managed_handlers() {
6363 let mut strategy = TimerOverrideStrategy {
6364 core: StrategyCore::new(StrategyConfig {
6365 strategy_id: Some(StrategyId::from("TEST-001")),
6366 ..Default::default()
6367 }),
6368 gtd_expiries: 0,
6369 market_exit_checks: 0,
6370 };
6371 let market_event = TimeEvent::new(
6372 strategy.core.market_exit_timer_name,
6373 UUID4::new(),
6374 UnixNanos::default(),
6375 UnixNanos::default(),
6376 );
6377
6378 route_time_event(&mut strategy, &market_event);
6379 DataActor::on_time_event(&mut strategy, &market_event).unwrap();
6380
6381 let client_order_id = ClientOrderId::from("O-001");
6382 let gtd_timer_name = Ustr::from("GTD-EXPIRY:O-001");
6383 strategy
6384 .core
6385 .gtd_timers
6386 .insert(client_order_id, gtd_timer_name);
6387 let gtd_event = TimeEvent::new(
6388 gtd_timer_name,
6389 UUID4::new(),
6390 UnixNanos::default(),
6391 UnixNanos::default(),
6392 );
6393
6394 route_time_event(&mut strategy, >d_event);
6395 DataActor::on_time_event(&mut strategy, >d_event).unwrap();
6396
6397 assert_eq!(strategy.market_exit_checks, 1);
6398 assert_eq!(strategy.gtd_expiries, 1);
6399 assert!(strategy.core.gtd_timers.contains_key(&client_order_id));
6400 assert_eq!(
6401 strategy.core.managed_time_event_last_id,
6402 Some(gtd_event.event_id)
6403 );
6404 }
6405
6406 #[rstest]
6407 fn test_route_time_event_ignores_unowned_market_exit_timer() {
6408 let mut strategy = TestStrategy::new(StrategyConfig {
6409 strategy_id: Some(StrategyId::from("TEST-001")),
6410 ..Default::default()
6411 });
6412 register_strategy(&mut strategy);
6413
6414 let order = make_accepted_limit_order("O-PRESENT");
6415 add_order_to_cache(&strategy, &order);
6416 strategy.core.cache_rc().borrow_mut().build_index();
6417 strategy.core.is_exiting = true;
6418 let event = TimeEvent::new(
6419 Ustr::from("MARKET_EXIT_CHECK:OTHER-001"),
6420 UUID4::new(),
6421 UnixNanos::default(),
6422 UnixNanos::default(),
6423 );
6424
6425 route_time_event(&mut strategy, &event);
6426
6427 assert!(strategy.core.is_exiting);
6428 assert_eq!(strategy.core.market_exit_attempts, 0);
6429 assert_eq!(strategy.core.managed_time_event_last_id, None);
6430 }
6431
6432 #[rstest]
6433 fn test_check_market_exit_finalizes_when_max_attempts_reached() {
6434 let config = StrategyConfig {
6435 strategy_id: Some(StrategyId::from("TEST-001")),
6436 order_id_tag: Some("001".to_string()),
6437 market_exit_max_attempts: 3,
6438 ..Default::default()
6439 };
6440 let mut strategy = TestStrategy::new(config);
6441 register_strategy(&mut strategy);
6442
6443 strategy.core.is_exiting = true;
6444 strategy.core.market_exit_attempts = 2; let event = TimeEvent::new(
6447 Ustr::from("MARKET_EXIT_CHECK:TEST-001"),
6448 UUID4::new(),
6449 UnixNanos::default(),
6450 UnixNanos::default(),
6451 );
6452 strategy.check_market_exit(event);
6453
6454 assert!(!strategy.core.is_exiting);
6456 assert_eq!(strategy.core.market_exit_attempts, 0);
6457 }
6458
6459 #[rstest]
6460 fn test_check_market_exit_finalizes_when_no_orders_or_positions() {
6461 let mut strategy = create_test_strategy();
6462 register_strategy(&mut strategy);
6463
6464 strategy.core.is_exiting = true;
6465
6466 let event = TimeEvent::new(
6467 Ustr::from("MARKET_EXIT_CHECK:TEST-001"),
6468 UUID4::new(),
6469 UnixNanos::default(),
6470 UnixNanos::default(),
6471 );
6472 strategy.check_market_exit(event);
6473
6474 assert!(!strategy.core.is_exiting);
6476 }
6477
6478 #[rstest]
6479 fn test_market_exit_timer_name_format() {
6480 let config = StrategyConfig {
6481 strategy_id: Some(StrategyId::from("MY-STRATEGY-001")),
6482 ..Default::default()
6483 };
6484 let strategy = TestStrategy::new(config);
6485
6486 assert_eq!(
6487 strategy.core.market_exit_timer_name,
6488 "MARKET_EXIT_CHECK:MY-STRATEGY-001"
6489 );
6490 }
6491
6492 #[rstest]
6493 fn test_reset_market_exit_state() {
6494 let mut strategy = create_test_strategy();
6495
6496 strategy.core.is_exiting = true;
6497 strategy.core.pending_stop = true;
6498 strategy.core.market_exit_attempts = 50;
6499
6500 strategy.core.reset_market_exit_state();
6501
6502 assert!(!strategy.core.is_exiting);
6503 assert!(!strategy.core.pending_stop);
6504 assert_eq!(strategy.core.market_exit_attempts, 0);
6505 }
6506
6507 #[rstest]
6508 fn test_cancel_market_exit_resets_state_without_hooks() {
6509 let config = StrategyConfig {
6510 strategy_id: Some(StrategyId::from("TEST-001")),
6511 order_id_tag: Some("001".to_string()),
6512 ..Default::default()
6513 };
6514 let mut strategy = MarketExitHookTrackingStrategy::new(config);
6515
6516 let trader_id = TraderId::from("TRADER-001");
6517 let clock = Rc::new(RefCell::new(TestClock::new()));
6518 let cache = Rc::new(RefCell::new(Cache::default()));
6519 let portfolio = Rc::new(RefCell::new(Portfolio::new(
6520 clock.clone(),
6521 cache.clone(),
6522 None,
6523 )));
6524 strategy
6525 .core
6526 .register(trader_id, clock, cache, portfolio)
6527 .unwrap();
6528
6529 strategy.core.is_exiting = true;
6531 strategy.core.pending_stop = true;
6532 strategy.core.market_exit_attempts = 50;
6533
6534 strategy.cancel_market_exit();
6536
6537 assert!(!strategy.core.is_exiting);
6539 assert!(!strategy.core.pending_stop);
6540 assert_eq!(strategy.core.market_exit_attempts, 0);
6541
6542 assert!(!strategy.on_market_exit_called);
6544 assert!(!strategy.post_market_exit_called);
6545 }
6546
6547 #[rstest]
6548 fn test_market_exit_returns_early_when_not_running() {
6549 let mut strategy = create_test_strategy();
6550 register_strategy(&mut strategy);
6551
6552 assert!(!strategy.is_running());
6554
6555 let result = strategy.market_exit();
6556
6557 assert!(result.is_ok());
6559 assert!(!strategy.core.is_exiting);
6560 }
6561
6562 #[rstest]
6563 fn test_stop_with_manage_stop_false_cleans_up_active_exit() {
6564 let config = StrategyConfig {
6565 strategy_id: Some(StrategyId::from("TEST-001")),
6566 order_id_tag: Some("001".to_string()),
6567 manage_stop: false,
6568 ..Default::default()
6569 };
6570 let mut strategy = TestStrategy::new(config);
6571 register_strategy(&mut strategy);
6572
6573 strategy.core.is_exiting = true;
6575 strategy.core.market_exit_attempts = 5;
6576
6577 let should_proceed = Strategy::stop(&mut strategy);
6579
6580 assert!(should_proceed);
6582 assert!(!strategy.core.is_exiting);
6583 assert_eq!(strategy.core.market_exit_attempts, 0);
6584 }
6585
6586 #[rstest]
6587 fn test_stop_with_manage_stop_true_defers_when_running() {
6588 let config = StrategyConfig {
6589 strategy_id: Some(StrategyId::from("TEST-001")),
6590 order_id_tag: Some("001".to_string()),
6591 manage_stop: true,
6592 ..Default::default()
6593 };
6594 let mut strategy = TestStrategy::new(config);
6595
6596 let trader_id = TraderId::from("TRADER-001");
6598 let clock = Rc::new(RefCell::new(TestClock::new()));
6599 clock
6600 .borrow_mut()
6601 .register_default_handler(TimeEventCallback::from(|_event: TimeEvent| {}));
6602 let cache = Rc::new(RefCell::new(Cache::default()));
6603 let portfolio = Rc::new(RefCell::new(Portfolio::new(
6604 clock.clone(),
6605 cache.clone(),
6606 None,
6607 )));
6608 strategy
6609 .core
6610 .register(trader_id, clock, cache, portfolio)
6611 .unwrap();
6612 strategy.initialize().unwrap();
6613 strategy.start().unwrap();
6614
6615 let should_proceed = Strategy::stop(&mut strategy);
6616
6617 assert!(!should_proceed);
6619 assert!(strategy.core.pending_stop);
6620 }
6621
6622 #[rstest]
6623 fn test_stop_with_manage_stop_true_returns_early_if_pending() {
6624 let config = StrategyConfig {
6625 strategy_id: Some(StrategyId::from("TEST-001")),
6626 order_id_tag: Some("001".to_string()),
6627 manage_stop: true,
6628 ..Default::default()
6629 };
6630 let mut strategy = TestStrategy::new(config);
6631 register_strategy(&mut strategy);
6632 start_strategy(&mut strategy);
6633 strategy.core.pending_stop = true;
6634
6635 let should_proceed = Strategy::stop(&mut strategy);
6637
6638 assert!(!should_proceed);
6640 assert!(strategy.core.pending_stop);
6641 }
6642
6643 #[rstest]
6644 fn test_stop_with_manage_stop_true_proceeds_when_not_running() {
6645 let config = StrategyConfig {
6646 strategy_id: Some(StrategyId::from("TEST-001")),
6647 order_id_tag: Some("001".to_string()),
6648 manage_stop: true,
6649 ..Default::default()
6650 };
6651 let mut strategy = TestStrategy::new(config);
6652 register_strategy(&mut strategy);
6653
6654 assert!(!strategy.is_running());
6656
6657 let should_proceed = Strategy::stop(&mut strategy);
6658
6659 assert!(should_proceed);
6661 }
6662
6663 #[rstest]
6664 fn test_finalize_market_exit_stops_strategy_when_pending() {
6665 let config = StrategyConfig {
6666 strategy_id: Some(StrategyId::from("TEST-001")),
6667 order_id_tag: Some("001".to_string()),
6668 ..Default::default()
6669 };
6670 let mut strategy = TestStrategy::new(config);
6671 register_strategy(&mut strategy);
6672 start_strategy(&mut strategy);
6673
6674 strategy.core.is_exiting = true;
6676 strategy.core.pending_stop = true;
6677
6678 strategy.finalize_market_exit();
6679
6680 assert_eq!(strategy.state(), ComponentState::Stopped);
6682 assert!(!strategy.core.is_exiting);
6683 assert!(!strategy.core.pending_stop);
6684 }
6685
6686 #[rstest]
6687 fn test_finalize_market_exit_stays_running_when_not_pending() {
6688 let config = StrategyConfig {
6689 strategy_id: Some(StrategyId::from("TEST-001")),
6690 order_id_tag: Some("001".to_string()),
6691 ..Default::default()
6692 };
6693 let mut strategy = TestStrategy::new(config);
6694 register_strategy(&mut strategy);
6695 start_strategy(&mut strategy);
6696
6697 strategy.core.is_exiting = true;
6699 strategy.core.pending_stop = false;
6700
6701 strategy.finalize_market_exit();
6702
6703 assert_eq!(strategy.state(), ComponentState::Running);
6705 assert!(!strategy.core.is_exiting);
6706 }
6707
6708 #[rstest]
6709 fn test_submit_order_denied_during_market_exit_when_not_reduce_only() {
6710 let mut strategy = create_test_strategy();
6711 register_strategy(&mut strategy);
6712 start_strategy(&mut strategy);
6713 strategy.core.is_exiting = true;
6714
6715 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
6716 get_typed_message_saving_handler(Some(Ustr::from("events.order.denied")));
6717 let order = OrderAny::Market(MarketOrder::new(
6718 TraderId::from("TRADER-001"),
6719 StrategyId::from("TEST-001"),
6720 InstrumentId::from("BTCUSDT.BINANCE"),
6721 ClientOrderId::from("O-20250208-0001"),
6722 OrderSide::Buy,
6723 Quantity::from(100_000),
6724 TimeInForce::Gtc,
6725 UUID4::new(),
6726 UnixNanos::default(),
6727 false, false,
6729 None,
6730 None,
6731 None,
6732 None,
6733 None,
6734 None,
6735 None,
6736 None,
6737 ));
6738 let topic = format!("events.order.{}", order.strategy_id());
6739 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
6740 let client_order_id = order.client_order_id();
6741 let result = strategy.submit_order(order.clone(), None, None, None);
6742
6743 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
6744
6745 assert!(result.is_ok());
6746 let cache = strategy.cache();
6747 let cached_order = cache.order(&client_order_id).unwrap();
6748 assert_eq!(cached_order.status(), OrderStatus::Denied);
6749
6750 let event_messages = event_messages.get_messages();
6751 assert_eq!(event_messages.len(), 2);
6752 assert_eq!(
6753 event_messages[0],
6754 OrderEventAny::Initialized(order.init_event().clone())
6755 );
6756 let OrderEventAny::Denied(denied) = &event_messages[1] else {
6757 panic!("expected OrderDenied event");
6758 };
6759 assert_eq!(denied.reason, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
6760 }
6761
6762 #[rstest]
6763 fn test_submit_order_list_denied_during_market_exit_publishes_init_then_denied_events() {
6764 let mut strategy = create_test_strategy();
6765 register_strategy(&mut strategy);
6766 start_strategy(&mut strategy);
6767 strategy.core.is_exiting = true;
6768
6769 let orders = vec![
6770 make_initialized_market_order("O-20250208-LIST-DENY-001"),
6771 make_initialized_market_order("O-20250208-LIST-DENY-002"),
6772 ];
6773 let client_order_id1 = orders[0].client_order_id();
6774 let client_order_id2 = orders[1].client_order_id();
6775 let cache_rc = strategy.core.cache_rc();
6776 let timeline = Rc::new(RefCell::new(Vec::new()));
6777 let event_messages = Rc::new(RefCell::new(Vec::new()));
6778
6779 let event_handler = {
6780 let event_messages = event_messages.clone();
6781 let timeline = timeline.clone();
6782 TypedHandler::from_with_id("events.order.list_denied", move |event: &OrderEventAny| {
6783 match event {
6784 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id1 => {
6785 assert!(cache_rc.borrow().order_exists(&client_order_id1));
6786 timeline.borrow_mut().push("init1");
6787 }
6788 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id2 => {
6789 assert!(cache_rc.borrow().order_exists(&client_order_id2));
6790 timeline.borrow_mut().push("init2");
6791 }
6792 OrderEventAny::Denied(e) if e.client_order_id == client_order_id1 => {
6793 assert_eq!(e.reason, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
6794 let cache = cache_rc.borrow();
6795 let cached_order = cache.order(&client_order_id1).unwrap();
6796 assert_eq!(cached_order.status(), OrderStatus::Denied);
6797 timeline.borrow_mut().push("denied1");
6798 }
6799 OrderEventAny::Denied(e) if e.client_order_id == client_order_id2 => {
6800 assert_eq!(e.reason, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
6801 let cache = cache_rc.borrow();
6802 let cached_order = cache.order(&client_order_id2).unwrap();
6803 assert_eq!(cached_order.status(), OrderStatus::Denied);
6804 timeline.borrow_mut().push("denied2");
6805 }
6806 _ => panic!("unexpected order event {event:?}"),
6807 }
6808 event_messages.borrow_mut().push(event.clone());
6809 })
6810 };
6811 let risk_handler = {
6812 let timeline = timeline.clone();
6813 TypedIntoHandler::from_with_id(
6814 "RiskEngine.queue_execute",
6815 move |_command: TradingCommand| {
6816 timeline.borrow_mut().push("command");
6817 },
6818 )
6819 };
6820 msgbus::register_trading_command_endpoint(
6821 MessagingSwitchboard::risk_engine_queue_execute(),
6822 risk_handler,
6823 );
6824
6825 let topic = format!("events.order.{}", orders[0].strategy_id());
6826 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
6827 let result = strategy.submit_order_list(orders.clone(), None, None, None);
6828
6829 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
6830
6831 assert!(result.is_ok());
6832
6833 let cache = strategy.cache();
6834 let cached_order1 = cache.order(&client_order_id1).unwrap();
6835 let cached_order2 = cache.order(&client_order_id2).unwrap();
6836 assert_eq!(cached_order1.status(), OrderStatus::Denied);
6837 assert_eq!(cached_order2.status(), OrderStatus::Denied);
6838
6839 let event_messages = event_messages.borrow();
6840 assert_eq!(event_messages.len(), 4);
6841 assert_eq!(
6842 event_messages[0],
6843 OrderEventAny::Initialized(orders[0].init_event().clone())
6844 );
6845 assert!(matches!(
6846 &event_messages[1],
6847 OrderEventAny::Denied(e)
6848 if e.client_order_id == client_order_id1
6849 && e.reason == Ustr::from("MARKET_EXIT_IN_PROGRESS")
6850 ));
6851 assert_eq!(
6852 event_messages[2],
6853 OrderEventAny::Initialized(orders[1].init_event().clone())
6854 );
6855 assert!(matches!(
6856 &event_messages[3],
6857 OrderEventAny::Denied(e)
6858 if e.client_order_id == client_order_id2
6859 && e.reason == Ustr::from("MARKET_EXIT_IN_PROGRESS")
6860 ));
6861 assert_eq!(
6862 timeline.borrow().as_slice(),
6863 &["init1", "denied1", "init2", "denied2"]
6864 );
6865 }
6866
6867 #[rstest]
6868 fn test_submit_order_list_market_exit_rejects_non_initialized_without_events() {
6869 let mut strategy = create_test_strategy();
6870 register_strategy(&mut strategy);
6871 start_strategy(&mut strategy);
6872 strategy.core.is_exiting = true;
6873
6874 let order = make_accepted_market_order("O-20250208-LIST-DENY-ACCEPTED");
6875 let topic = format!("events.order.{}", order.strategy_id());
6876 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
6877 get_typed_message_saving_handler(Some(Ustr::from("events.order.list_invalid")));
6878
6879 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
6880 let result = strategy.submit_order_list(vec![order], None, None, None);
6881
6882 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
6883
6884 assert!(result.is_err());
6885 assert!(
6886 result
6887 .unwrap_err()
6888 .to_string()
6889 .contains("expected INITIALIZED")
6890 );
6891 assert!(event_messages.get_messages().is_empty());
6892 }
6893
6894 #[rstest]
6895 fn test_submit_order_list_rejects_mixed_venues_with_friendly_error() {
6896 let mut strategy = create_test_strategy();
6897 register_strategy(&mut strategy);
6898 start_strategy(&mut strategy);
6899
6900 let binance_order = make_initialized_market_order("O-MIXED-VENUE-001");
6901 let bybit_order = OrderAny::Market(MarketOrder::new(
6902 TraderId::from("TRADER-001"),
6903 StrategyId::from("TEST-001"),
6904 InstrumentId::from("BTCUSDT.BYBIT"),
6905 ClientOrderId::from("O-MIXED-VENUE-002"),
6906 OrderSide::Buy,
6907 Quantity::from(100_000),
6908 TimeInForce::Gtc,
6909 UUID4::new(),
6910 UnixNanos::default(),
6911 false,
6912 false,
6913 None,
6914 None,
6915 None,
6916 None,
6917 None,
6918 None,
6919 None,
6920 None,
6921 ));
6922
6923 let result = strategy.submit_order_list(vec![binance_order, bybit_order], None, None, None);
6924
6925 let err = result.unwrap_err();
6926 let msg = err.to_string();
6927 assert!(
6928 msg.contains("OrderList denied: orders must share the same venue"),
6929 "unexpected error: {msg}",
6930 );
6931 assert!(msg.contains("BINANCE"), "expected BINANCE in error: {msg}");
6932 assert!(msg.contains("BYBIT"), "expected BYBIT in error: {msg}");
6933 }
6934
6935 #[rstest]
6936 fn test_submit_order_allowed_during_market_exit_when_reduce_only() {
6937 let mut strategy = create_test_strategy();
6938 register_strategy(&mut strategy);
6939 start_strategy(&mut strategy);
6940 strategy.core.is_exiting = true;
6941
6942 let order = OrderAny::Market(MarketOrder::new(
6943 TraderId::from("TRADER-001"),
6944 StrategyId::from("TEST-001"),
6945 InstrumentId::from("BTCUSDT.BINANCE"),
6946 ClientOrderId::from("O-20250208-0001"),
6947 OrderSide::Buy,
6948 Quantity::from(100_000),
6949 TimeInForce::Gtc,
6950 UUID4::new(),
6951 UnixNanos::default(),
6952 true, false,
6954 None,
6955 None,
6956 None,
6957 None,
6958 None,
6959 None,
6960 None,
6961 None,
6962 ));
6963 let client_order_id = order.client_order_id();
6964 let result = strategy.submit_order(order, None, None, None);
6965
6966 assert!(result.is_ok());
6967 let cache = strategy.cache();
6968 let cached_order = cache.order(&client_order_id).unwrap();
6969 assert_ne!(cached_order.status(), OrderStatus::Denied);
6970 }
6971
6972 #[rstest]
6973 fn test_submit_order_allowed_during_market_exit_when_tagged() {
6974 let mut strategy = create_test_strategy();
6975 register_strategy(&mut strategy);
6976 start_strategy(&mut strategy);
6977 strategy.core.is_exiting = true;
6978
6979 let order = OrderAny::Market(MarketOrder::new(
6980 TraderId::from("TRADER-001"),
6981 StrategyId::from("TEST-001"),
6982 InstrumentId::from("BTCUSDT.BINANCE"),
6983 ClientOrderId::from("O-20250208-0002"),
6984 OrderSide::Buy,
6985 Quantity::from(100_000),
6986 TimeInForce::Gtc,
6987 UUID4::new(),
6988 UnixNanos::default(),
6989 false, false,
6991 None,
6992 None,
6993 None,
6994 None,
6995 None,
6996 None,
6997 None,
6998 Some(vec![Ustr::from("MARKET_EXIT")]),
6999 ));
7000 let client_order_id = order.client_order_id();
7001 let result = strategy.submit_order(order, None, None, None);
7002
7003 assert!(result.is_ok());
7004 let cache = strategy.cache();
7005 let cached_order = cache.order(&client_order_id).unwrap();
7006 assert_ne!(cached_order.status(), OrderStatus::Denied);
7007 }
7008
7009 #[derive(Debug)]
7010 struct MacroTestSimple {
7011 core: StrategyCore,
7012 }
7013
7014 nautilus_strategy!(MacroTestSimple);
7015
7016 impl DataActor for MacroTestSimple {}
7017
7018 #[derive(Debug)]
7019 struct MacroTestWithHooks {
7020 core: StrategyCore,
7021 }
7022
7023 nautilus_strategy!(MacroTestWithHooks, {
7024 fn on_order_rejected(&mut self, _event: OrderRejected) {}
7025 });
7026
7027 impl DataActor for MacroTestWithHooks {}
7028
7029 #[derive(Debug)]
7030 struct MacroTestCustomField {
7031 inner: StrategyCore,
7032 }
7033
7034 nautilus_strategy!(MacroTestCustomField, inner, {
7035 fn external_order_instrument_ids(&self) -> Option<Vec<InstrumentId>> {
7036 None
7037 }
7038 });
7039
7040 impl DataActor for MacroTestCustomField {}
7041
7042 #[rstest]
7043 fn test_strategy_behavior_does_not_require_native_core_access() {
7044 fn assert_strategy<T: Strategy + DataActor>() {}
7045
7046 assert_strategy::<CoreFreeStrategy>();
7047
7048 let mut strategy = CoreFreeStrategy { started: false };
7049 DataActor::on_start(&mut strategy).unwrap();
7050
7051 assert!(strategy.started);
7052 }
7053
7054 #[rstest]
7055 fn test_nautilus_strategy_macro_forms() {
7056 let config = StrategyConfig {
7057 strategy_id: Some(StrategyId::from("MACRO-001")),
7058 order_id_tag: Some("001".to_string()),
7059 ..Default::default()
7060 };
7061
7062 let simple = MacroTestSimple {
7063 core: StrategyCore::new(config.clone()),
7064 };
7065 assert_eq!(simple.strategy_id(), config.strategy_id);
7066 assert_eq!(simple.config().order_id_tag, config.order_id_tag);
7067 assert_eq!(simple.actor_id(), ActorId::from("MACRO-001"));
7068
7069 let hooks = MacroTestWithHooks {
7070 core: StrategyCore::new(config.clone()),
7071 };
7072 assert_eq!(hooks.strategy_id(), config.strategy_id);
7073 assert_eq!(hooks.config().order_id_tag, config.order_id_tag);
7074 assert_eq!(hooks.actor_id(), ActorId::from("MACRO-001"));
7075
7076 let custom = MacroTestCustomField {
7077 inner: StrategyCore::new(config.clone()),
7078 };
7079 assert_eq!(custom.strategy_id(), config.strategy_id);
7080 assert_eq!(custom.config().order_id_tag, config.order_id_tag);
7081 assert_eq!(custom.actor_id(), ActorId::from("MACRO-001"));
7082 assert!(custom.external_order_instrument_ids().is_none());
7083 }
7084}