1pub mod api;
17pub mod config;
18pub mod core;
19
20pub use core::{StrategyCore, StrategyNative};
21use std::panic::{AssertUnwindSafe, catch_unwind};
22
23use ahash::AHashSet;
24pub use api::{OrderApi, PortfolioApi};
25pub use config::{ImportableStrategyConfig, StrategyConfig};
26use nautilus_common::{
27 actor::DataActor,
28 component::Component,
29 enums::ComponentState,
30 logging::{CMD, EVT, RECV, SEND},
31 messages::execution::{
32 BatchCancelOrders, BatchModifyOrders, CancelAllOrders, CancelOrder, ModifyOrder,
33 QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList, TradingCommand,
34 },
35 msgbus::{self, MessagingSwitchboard},
36 timer::TimeEvent,
37};
38use nautilus_core::{Params, UUID4};
39use nautilus_model::{
40 enums::{OrderSide, OrderStatus, PositionSide, TimeInForce, TriggerType},
41 events::{
42 OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied, OrderEmulated,
43 OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled, OrderInitialized,
44 OrderModifyRejected, OrderPendingCancel, OrderPendingUpdate, OrderRejected, OrderReleased,
45 OrderSubmitted, OrderTriggered, OrderUpdated, PositionChanged, PositionClosed,
46 PositionEvent, PositionOpened,
47 },
48 identifiers::{
49 AccountId, ClientId, ClientOrderId, ExecAlgorithmId, InstrumentId, PositionId, StrategyId,
50 TraderId,
51 },
52 orders::{
53 LIMIT_ORDER_TYPES, Order, OrderAny, OrderCore, OrderError, OrderList, STOP_ORDER_TYPES,
54 },
55 position::Position,
56 types::{Price, Quantity},
57};
58use ustr::Ustr;
59
60pub type BatchModifyOrder = (
62 ClientOrderId,
63 Option<Quantity>,
64 Option<Price>,
65 Option<Price>,
66);
67
68pub trait Strategy: DataActor {
106 fn external_order_claims(&self) -> Option<Vec<InstrumentId>> {
111 None
112 }
113
114 fn strategy_id(&self) -> Option<StrategyId>
116 where
117 Self: StrategyNative,
118 {
119 StrategyNative::strategy_core(self).strategy_id()
120 }
121
122 fn order(&self) -> OrderApi<'_>
124 where
125 Self: StrategyNative,
126 {
127 StrategyNative::strategy_core(self).order()
128 }
129
130 fn portfolio(&self) -> PortfolioApi<'_>
132 where
133 Self: StrategyNative,
134 {
135 StrategyNative::strategy_core(self).portfolio_api()
136 }
137
138 fn submit_order(
144 &mut self,
145 order: OrderAny,
146 position_id: Option<PositionId>,
147 client_id: Option<ClientId>,
148 params: Option<Params>,
149 ) -> anyhow::Result<()>
150 where
151 Self: StrategyNative,
152 {
153 let core = StrategyNative::strategy_core_mut(self);
154
155 let trader_id = registered_trader_id(core)?;
156 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
157 let ts_init = core.clock_mut().timestamp_ns();
158
159 if order.status() != OrderStatus::Initialized {
160 anyhow::bail!(
161 "Order denied: invalid status for {}, expected INITIALIZED",
162 order.client_order_id()
163 );
164 }
165
166 let market_exit_tag = core.market_exit_tag;
167 let is_market_exit_order = order
168 .tags()
169 .is_some_and(|tags| tags.contains(&market_exit_tag));
170 let should_deny_for_market_exit =
171 core.is_exiting && !order.is_reduce_only() && !is_market_exit_order;
172
173 if should_deny_for_market_exit {
174 self.deny_order(&order, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
175 return Ok(());
176 }
177
178 let core = StrategyNative::strategy_core_mut(self);
179 let params = params.filter(|params| !params.is_empty());
180
181 {
182 let cache_rc = core.cache_rc();
183 let mut cache = cache_rc.try_borrow_mut().map_err(|_| {
184 anyhow::anyhow!(
185 "Cannot submit order {}: cache is currently borrowed",
186 order.client_order_id()
187 )
188 })?;
189 cache.add_order(order.clone(), position_id, client_id, true)?;
190 }
191
192 publish_order_initialized(&order);
193
194 let command = SubmitOrder::new(
195 trader_id,
196 client_id,
197 strategy_id,
198 order.instrument_id(),
199 order.client_order_id(),
200 order.init_event().clone(),
201 order.exec_algorithm_id(),
202 position_id,
203 params,
204 UUID4::new(),
205 ts_init,
206 None, );
208
209 if matches!(order.emulation_trigger(), Some(trigger) if trigger != TriggerType::NoTrigger) {
210 send_emulator_command(TradingCommand::SubmitOrder(command));
211 } else if let Some(exec_algorithm_id) = order.exec_algorithm_id() {
212 send_algo_command(command, exec_algorithm_id);
213 } else {
214 send_risk_command(TradingCommand::SubmitOrder(command));
215 }
216
217 self.set_gtd_expiry(&order)?;
218 Ok(())
219 }
220
221 fn submit_order_list(
228 &mut self,
229 mut orders: Vec<OrderAny>,
230 position_id: Option<PositionId>,
231 client_id: Option<ClientId>,
232 params: Option<Params>,
233 ) -> anyhow::Result<()>
234 where
235 Self: StrategyNative,
236 {
237 if orders.is_empty() {
238 log::error!("OrderList denied: no orders to submit");
239 anyhow::bail!("OrderList denied: no orders to submit");
240 }
241
242 for order in &orders {
243 if order.status() != OrderStatus::Initialized {
244 anyhow::bail!(
245 "Order in list denied: invalid status for {}, expected INITIALIZED",
246 order.client_order_id()
247 );
248 }
249 }
250
251 let first_venue = orders[0].instrument_id().venue;
252 for order in &orders {
253 if order.instrument_id().venue != first_venue {
254 anyhow::bail!(
255 "OrderList denied: orders must share the same venue; \
256 expected {first_venue}, found {} on {}",
257 order.instrument_id().venue,
258 order.client_order_id(),
259 );
260 }
261 }
262
263 let should_deny = {
264 let core = StrategyNative::strategy_core_mut(self);
265 let tag = core.market_exit_tag;
266 core.is_exiting
267 && orders.iter().any(|o| {
268 !o.is_reduce_only() && !o.tags().is_some_and(|tags| tags.contains(&tag))
269 })
270 };
271
272 if should_deny {
273 self.deny_order_list(&orders, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
274 return Ok(());
275 }
276
277 let core = StrategyNative::strategy_core_mut(self);
278
279 let trader_id = registered_trader_id(core)?;
280 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
281 let ts_init = core.clock_mut().timestamp_ns();
282
283 let order_list = if orders.first().is_some_and(|o| o.order_list_id().is_some()) {
285 OrderList::from_orders(&orders, ts_init)
286 } else {
287 core.order_factory().create_list(&mut orders, ts_init)
288 };
289
290 if let Err(e) = order_list.validate() {
291 log::error!("OrderList denied: {e}");
292 anyhow::bail!("OrderList denied: {e}");
293 }
294
295 {
296 let cache_rc = core.cache_rc();
297 let mut cache = cache_rc.try_borrow_mut().map_err(|_| {
298 anyhow::anyhow!(
299 "Cannot submit order list {}: cache is currently borrowed",
300 order_list.id
301 )
302 })?;
303
304 if cache.order_list_exists(&order_list.id) {
305 anyhow::bail!("OrderList denied: duplicate {}", order_list.id);
306 }
307
308 for order in &orders {
309 if cache.order_exists(&order.client_order_id()) {
310 anyhow::bail!(
311 "Order in list denied: duplicate {}",
312 order.client_order_id()
313 );
314 }
315 }
316
317 cache.add_order_list(order_list.clone())?;
318 for order in &orders {
319 cache.add_order(order.clone(), position_id, client_id, true)?;
320 }
321 }
322
323 for order in &orders {
324 publish_order_initialized(order);
325 }
326
327 let params = params.filter(|params| !params.is_empty());
328
329 let first_order = orders.first();
330 let order_inits: Vec<_> = orders.iter().map(|o| o.init_event().clone()).collect();
331 let exec_algorithm_id = first_order.and_then(|o| o.exec_algorithm_id());
332
333 let command = SubmitOrderList::new(
334 trader_id,
335 client_id,
336 strategy_id,
337 order_list,
338 order_inits,
339 exec_algorithm_id,
340 position_id,
341 params,
342 UUID4::new(),
343 ts_init,
344 None, );
346
347 let has_emulated_order = orders.iter().any(|o| {
348 matches!(o.emulation_trigger(), Some(trigger) if trigger != TriggerType::NoTrigger)
349 || o.is_emulated()
350 });
351
352 if has_emulated_order {
353 send_emulator_command(TradingCommand::SubmitOrderList(command));
354 } else if let Some(algo_id) = exec_algorithm_id {
355 let endpoint = format!("{algo_id}.execute");
356 msgbus::send_any(endpoint.into(), &TradingCommand::SubmitOrderList(command));
357 } else {
358 send_risk_command(TradingCommand::SubmitOrderList(command));
359 }
360
361 for order in &orders {
362 self.set_gtd_expiry(order)?;
363 }
364
365 Ok(())
366 }
367
368 fn modify_order(
374 &mut self,
375 client_order_id: ClientOrderId,
376 quantity: Option<Quantity>,
377 price: Option<Price>,
378 trigger_price: Option<Price>,
379 client_id: Option<ClientId>,
380 params: Option<Params>,
381 ) -> anyhow::Result<()>
382 where
383 Self: StrategyNative,
384 {
385 let (trader_id, strategy_id) = {
386 let core = StrategyNative::strategy_core_mut(self);
387 (
388 registered_trader_id(core)?,
389 StrategyId::from(core.actor_id().inner().as_str()),
390 )
391 };
392
393 let params = params.filter(|params| !params.is_empty());
394
395 let order = StrategyNative::strategy_core_mut(self)
397 .cache_rc()
398 .borrow()
399 .try_order_owned(&client_order_id)
400 .map_err(|e| anyhow::anyhow!("Cannot modify order: {e}"))?;
401
402 let mut updating = false;
403
404 if quantity.is_some_and(|q| q != order.quantity()) {
405 updating = true;
406 }
407
408 if let Some(price) = price {
409 if !LIMIT_ORDER_TYPES.contains(&order.order_type()) {
410 anyhow::bail!("{} orders do not have a LIMIT price", order.order_type());
411 }
412
413 if Some(price) != order.price() {
414 updating = true;
415 }
416 }
417
418 if let Some(trigger_price) = trigger_price {
419 if !STOP_ORDER_TYPES.contains(&order.order_type()) {
420 anyhow::bail!(
421 "{} orders do not have a STOP trigger price",
422 order.order_type()
423 );
424 }
425
426 if Some(trigger_price) != order.trigger_price() {
427 updating = true;
428 }
429 }
430
431 if !updating {
432 log::error!(
433 "Cannot create command ModifyOrder: quantity, price, and trigger were either None \
434 or the same as existing values"
435 );
436 return Ok(());
437 }
438
439 if order.is_closed() || order.is_pending_cancel() {
440 log::warn!(
441 "Cannot create command ModifyOrder: state is {:?}, {order:?}",
442 order.status()
443 );
444 return Ok(());
445 }
446
447 if !self.mark_order_pending_update(&order)? {
448 return Ok(());
449 }
450
451 let command = ModifyOrder::new(
452 trader_id,
453 client_id,
454 strategy_id,
455 order.instrument_id(),
456 order.client_order_id(),
457 order.venue_order_id(),
458 quantity,
459 price,
460 trigger_price,
461 UUID4::new(),
462 StrategyNative::strategy_core_mut(self)
463 .clock_mut()
464 .timestamp_ns(),
465 params,
466 None, );
468
469 if order.is_emulated() {
470 send_emulator_command(TradingCommand::ModifyOrder(command));
471 } else {
472 send_risk_command(TradingCommand::ModifyOrder(command));
473 }
474 Ok(())
475 }
476
477 fn modify_orders(
486 &mut self,
487 updates: Vec<BatchModifyOrder>,
488 client_id: Option<ClientId>,
489 params: Option<Params>,
490 ) -> anyhow::Result<()>
491 where
492 Self: StrategyNative,
493 {
494 if updates.is_empty() {
495 anyhow::bail!("Cannot batch modify empty order list");
496 }
497
498 let (trader_id, strategy_id, ts_init) = {
499 let core = StrategyNative::strategy_core_mut(self);
500 (
501 registered_trader_id(core)?,
502 StrategyId::from(core.actor_id().inner().as_str()),
503 core.clock_mut().timestamp_ns(),
504 )
505 };
506
507 let orders: Vec<OrderAny> = {
508 let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
509 let cache = cache_rc.borrow();
510 updates
511 .iter()
512 .map(|(client_order_id, _, _, _)| {
513 cache
514 .try_order_owned(client_order_id)
515 .map_err(|e| anyhow::anyhow!("Cannot modify order: {e}"))
516 })
517 .collect::<Result<_, _>>()?
518 };
519
520 let instrument_id = orders[0].instrument_id();
521
522 for (order, (_, quantity, price, trigger_price)) in orders.iter().zip(updates.iter()) {
523 if order.instrument_id() != instrument_id {
524 anyhow::bail!(
525 "Cannot batch modify orders for different instruments: {} vs {}",
526 instrument_id,
527 order.instrument_id()
528 );
529 }
530
531 if order.is_emulated() || order.is_active_local() {
532 anyhow::bail!("Cannot include emulated or local orders in batch modify");
533 }
534
535 let mut updating = false;
536
537 if quantity.is_some_and(|q| q != order.quantity()) {
538 updating = true;
539 }
540
541 if let Some(price) = price {
542 if !LIMIT_ORDER_TYPES.contains(&order.order_type()) {
543 anyhow::bail!("{} orders do not have a LIMIT price", order.order_type());
544 }
545
546 if Some(*price) != order.price() {
547 updating = true;
548 }
549 }
550
551 if let Some(trigger_price) = trigger_price {
552 if !STOP_ORDER_TYPES.contains(&order.order_type()) {
553 anyhow::bail!(
554 "{} orders do not have a STOP trigger price",
555 order.order_type()
556 );
557 }
558
559 if Some(*trigger_price) != order.trigger_price() {
560 updating = true;
561 }
562 }
563
564 if !updating {
565 anyhow::bail!(
566 "Cannot create command BatchModifyOrders: quantity, price, and trigger were \
567 either None or the same as existing values for {}",
568 order.client_order_id()
569 );
570 }
571
572 if order.is_closed() || order.is_pending_cancel() {
573 anyhow::bail!(
574 "Cannot create command BatchModifyOrders: state is {:?}, {order:?}",
575 order.status()
576 );
577 }
578 }
579
580 let params = params.filter(|params| !params.is_empty());
581 let mut modifies = Vec::with_capacity(orders.len());
582
583 for (order, (_, quantity, price, trigger_price)) in orders.into_iter().zip(updates) {
584 if !self.mark_order_pending_update(&order)? {
585 continue;
586 }
587
588 modifies.push(ModifyOrder::new(
589 trader_id,
590 client_id,
591 strategy_id,
592 instrument_id,
593 order.client_order_id(),
594 order.venue_order_id(),
595 quantity,
596 price,
597 trigger_price,
598 UUID4::new(),
599 ts_init,
600 params.clone(),
601 None, ));
603 }
604
605 if modifies.is_empty() {
606 log::warn!("Cannot send `BatchModifyOrders`, no valid modify commands");
607 return Ok(());
608 }
609
610 let command = BatchModifyOrders::new(
611 trader_id,
612 client_id,
613 strategy_id,
614 instrument_id,
615 modifies,
616 UUID4::new(),
617 ts_init,
618 params,
619 None, );
621
622 send_risk_command(TradingCommand::ModifyOrders(command));
623 Ok(())
624 }
625
626 fn cancel_order(
632 &mut self,
633 client_order_id: ClientOrderId,
634 client_id: Option<ClientId>,
635 params: Option<Params>,
636 ) -> anyhow::Result<()>
637 where
638 Self: StrategyNative,
639 {
640 let (trader_id, strategy_id, ts_init) = {
641 let core = StrategyNative::strategy_core_mut(self);
642 (
643 registered_trader_id(core)?,
644 StrategyId::from(core.actor_id().inner().as_str()),
645 core.clock_mut().timestamp_ns(),
646 )
647 };
648
649 let params = params.filter(|params| !params.is_empty());
650
651 let order = StrategyNative::strategy_core_mut(self)
655 .cache_rc()
656 .borrow()
657 .try_order_owned(&client_order_id)
658 .map_err(|e| anyhow::anyhow!("Cannot cancel order: {e}"))?;
659
660 if !self.mark_order_pending_cancel(&order)? {
661 return Ok(());
662 }
663
664 let command = CancelOrder::new(
665 trader_id,
666 client_id,
667 strategy_id,
668 order.instrument_id(),
669 order.client_order_id(),
670 order.venue_order_id(),
671 UUID4::new(),
672 ts_init,
673 params,
674 None, );
676
677 if matches!(order.emulation_trigger(), Some(trigger) if trigger != TriggerType::NoTrigger)
678 || order.is_emulated()
679 {
680 send_emulator_command(TradingCommand::CancelOrder(command));
681 } else if let Some(algo_id) = order
682 .exec_algorithm_id()
683 .filter(|_| order.is_active_local())
684 {
685 let endpoint = format!("{algo_id}.execute");
686 msgbus::send_any(endpoint.into(), &TradingCommand::CancelOrder(command));
687 } else {
688 send_exec_command(TradingCommand::CancelOrder(command));
689 }
690
691 if StrategyNative::strategy_core(self).config.manage_gtd_expiry
692 && order.time_in_force() == TimeInForce::Gtd
693 && self.has_gtd_expiry_timer(&order.client_order_id())
694 {
695 self.cancel_gtd_expiry(&order.client_order_id());
696 }
697
698 Ok(())
699 }
700
701 fn cancel_orders(
708 &mut self,
709 client_order_ids: Vec<ClientOrderId>,
710 client_id: Option<ClientId>,
711 params: Option<Params>,
712 ) -> anyhow::Result<()>
713 where
714 Self: StrategyNative,
715 {
716 if client_order_ids.is_empty() {
717 anyhow::bail!("Cannot batch cancel empty order list");
718 }
719
720 let (trader_id, strategy_id, ts_init) = {
721 let core = StrategyNative::strategy_core_mut(self);
722 (
723 registered_trader_id(core)?,
724 StrategyId::from(core.actor_id().inner().as_str()),
725 core.clock_mut().timestamp_ns(),
726 )
727 };
728
729 let orders: Vec<OrderAny> = {
731 let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
732 let cache = cache_rc.borrow();
733 client_order_ids
734 .iter()
735 .map(|id| {
736 cache
737 .try_order_owned(id)
738 .map_err(|e| anyhow::anyhow!("Cannot cancel order: {e}"))
739 })
740 .collect::<Result<_, _>>()?
741 };
742
743 let instrument_id = orders[0].instrument_id();
744
745 for order in &orders {
746 if order.instrument_id() != instrument_id {
747 anyhow::bail!(
748 "Cannot batch cancel orders for different instruments: {} vs {}",
749 instrument_id,
750 order.instrument_id()
751 );
752 }
753
754 if order.is_emulated() || order.is_active_local() {
755 anyhow::bail!("Cannot include emulated or local orders in batch cancel");
756 }
757 }
758
759 let mut cancels = Vec::with_capacity(orders.len());
760
761 for order in orders {
762 if !self.mark_order_pending_cancel(&order)? {
763 continue;
764 }
765
766 cancels.push(CancelOrder::new(
767 trader_id,
768 client_id,
769 strategy_id,
770 instrument_id,
771 order.client_order_id(),
772 order.venue_order_id(),
773 UUID4::new(),
774 ts_init,
775 params.clone(),
776 None, ));
778 }
779
780 if cancels.is_empty() {
781 log::warn!("Cannot send `BatchCancelOrders`, no valid cancel commands");
782 return Ok(());
783 }
784
785 let command = BatchCancelOrders::new(
786 trader_id,
787 client_id,
788 strategy_id,
789 instrument_id,
790 cancels,
791 UUID4::new(),
792 ts_init,
793 params,
794 None, );
796
797 send_exec_command(TradingCommand::CancelOrders(command));
798 Ok(())
799 }
800
801 fn mark_order_pending_update(&mut self, order: &OrderAny) -> anyhow::Result<bool>
807 where
808 Self: StrategyNative,
809 {
810 if order.is_active_local() {
811 return Ok(true);
812 }
813
814 let strategy_id = order.strategy_id();
815 required_account_id(order, "pending update")?;
816 let event = OrderEventAny::PendingUpdate(self.generate_order_pending_update(order));
817
818 {
819 let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
820 let mut cache = cache_rc.borrow_mut();
821 match cache.update_order(&event) {
822 Ok(_) => {}
823 Err(e)
824 if matches!(
825 e.downcast_ref::<OrderError>(),
826 Some(OrderError::InvalidStateTransition)
827 ) =>
828 {
829 log::warn!("InvalidStateTrigger: {e}, did not apply pending update event");
830 return Ok(false);
831 }
832 Err(e) => return Err(e),
833 }
834 }
835
836 let topic = format!("events.order.{strategy_id}");
837 msgbus::publish_order_event(topic.into(), &event);
838 msgbus::publish_order_event(
839 msgbus::switchboard::get_order_pending_update_topic(order.instrument_id()),
840 &event,
841 );
842
843 Ok(true)
844 }
845
846 fn mark_order_pending_cancel(&mut self, order: &OrderAny) -> anyhow::Result<bool>
852 where
853 Self: StrategyNative,
854 {
855 if order.is_closed() || order.is_pending_cancel() {
856 log::warn!(
857 "Cannot cancel order: state is {:?}, {order:?}",
858 order.status()
859 );
860 return Ok(false);
861 }
862
863 if order.is_active_local() {
864 return Ok(true);
865 }
866
867 let strategy_id = order.strategy_id();
868 required_account_id(order, "pending cancel")?;
869 let event = OrderEventAny::PendingCancel(self.generate_order_pending_cancel(order));
870
871 {
872 let cache_rc = StrategyNative::strategy_core_mut(self).cache_rc();
873 let mut cache = cache_rc.borrow_mut();
874 match cache.update_order(&event) {
875 Ok(_) => {}
876 Err(e)
877 if matches!(
878 e.downcast_ref::<OrderError>(),
879 Some(OrderError::InvalidStateTransition)
880 ) =>
881 {
882 log::warn!("InvalidStateTrigger: {e}, did not apply pending cancel event");
883 return Ok(false);
884 }
885 Err(e) => return Err(e),
886 }
887 cache.update_order_pending_cancel_local(order);
888 }
889
890 let topic = format!("events.order.{strategy_id}");
891 msgbus::publish_order_event(topic.into(), &event);
892 msgbus::publish_order_event(
893 msgbus::switchboard::get_order_pending_cancel_topic(order.instrument_id()),
894 &event,
895 );
896
897 Ok(true)
898 }
899
900 fn generate_order_pending_update(&mut self, order: &OrderAny) -> OrderPendingUpdate
902 where
903 Self: StrategyNative,
904 {
905 let ts_now = StrategyNative::strategy_core_mut(self)
906 .clock_mut()
907 .timestamp_ns();
908
909 OrderPendingUpdate::new(
910 order.trader_id(),
911 order.strategy_id(),
912 order.instrument_id(),
913 order.client_order_id(),
914 order.account_id(),
915 UUID4::new(),
916 ts_now,
917 ts_now,
918 false,
919 order.venue_order_id(),
920 )
921 }
922
923 fn generate_order_pending_cancel(&mut self, order: &OrderAny) -> OrderPendingCancel
925 where
926 Self: StrategyNative,
927 {
928 let ts_now = StrategyNative::strategy_core_mut(self)
929 .clock_mut()
930 .timestamp_ns();
931
932 OrderPendingCancel::new(
933 order.trader_id(),
934 order.strategy_id(),
935 order.instrument_id(),
936 order.client_order_id(),
937 order.account_id(),
938 UUID4::new(),
939 ts_now,
940 ts_now,
941 false,
942 order.venue_order_id(),
943 )
944 }
945
946 fn cancel_all_orders(
952 &mut self,
953 instrument_id: InstrumentId,
954 order_side: Option<OrderSide>,
955 client_id: Option<ClientId>,
956 params: Option<Params>,
957 ) -> anyhow::Result<()>
958 where
959 Self: StrategyNative,
960 {
961 let params = params.filter(|params| !params.is_empty());
962 let core = StrategyNative::strategy_core_mut(self);
963
964 let trader_id = registered_trader_id(core)?;
965 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
966 let ts_init = core.clock_mut().timestamp_ns();
967 let cache = core.cache_ref();
968
969 let open_count = cache.orders_open_count(
970 None,
971 Some(&instrument_id),
972 Some(&strategy_id),
973 None,
974 order_side,
975 );
976
977 let emulated_count = cache.orders_emulated_count(
978 None,
979 Some(&instrument_id),
980 Some(&strategy_id),
981 None,
982 order_side,
983 );
984
985 let inflight_count = cache.orders_inflight_count(
986 None,
987 Some(&instrument_id),
988 Some(&strategy_id),
989 None,
990 order_side,
991 );
992
993 let mut exec_algorithm_ids: Vec<_> = cache.exec_algorithm_ids().into_iter().collect();
997 exec_algorithm_ids.sort();
998 let mut algo_orders: Vec<OrderAny> = Vec::new();
999
1000 for algo_id in &exec_algorithm_ids {
1001 algo_orders.extend(
1002 cache
1003 .orders_for_exec_algorithm(
1004 algo_id,
1005 None,
1006 Some(&instrument_id),
1007 Some(&strategy_id),
1008 None,
1009 order_side,
1010 )
1011 .into_iter()
1012 .map(|o| o.clone()),
1013 );
1014 }
1015
1016 let algo_count = algo_orders.len();
1017
1018 drop(cache);
1019
1020 if open_count == 0 && emulated_count == 0 && inflight_count == 0 && algo_count == 0 {
1021 let side_str = order_side.map(|s| format!(" {s}")).unwrap_or_default();
1022 log::info!("No {instrument_id} open, emulated, or inflight{side_str} orders to cancel");
1023 return Ok(());
1024 }
1025
1026 let side_str = order_side.map(|s| format!(" {s}")).unwrap_or_default();
1027
1028 if open_count > 0 {
1029 log::info!(
1030 "Canceling {open_count} open{side_str} {instrument_id} order{}",
1031 if open_count == 1 { "" } else { "s" }
1032 );
1033 }
1034
1035 if emulated_count > 0 {
1036 log::info!(
1037 "Canceling {emulated_count} emulated{side_str} {instrument_id} order{}",
1038 if emulated_count == 1 { "" } else { "s" }
1039 );
1040 }
1041
1042 if inflight_count > 0 {
1043 log::info!(
1044 "Canceling {inflight_count} inflight{side_str} {instrument_id} order{}",
1045 if inflight_count == 1 { "" } else { "s" }
1046 );
1047 }
1048
1049 if open_count > 0 || inflight_count > 0 {
1050 let command = CancelAllOrders::new(
1051 trader_id,
1052 client_id,
1053 strategy_id,
1054 instrument_id,
1055 order_side.unwrap_or(OrderSide::NoOrderSide),
1056 UUID4::new(),
1057 ts_init,
1058 params.clone(),
1059 None, );
1061
1062 send_exec_command(TradingCommand::CancelAllOrders(command));
1063 }
1064
1065 if emulated_count > 0 {
1066 let command = CancelAllOrders::new(
1067 trader_id,
1068 client_id,
1069 strategy_id,
1070 instrument_id,
1071 order_side.unwrap_or(OrderSide::NoOrderSide),
1072 UUID4::new(),
1073 ts_init,
1074 params,
1075 None, );
1077
1078 send_emulator_command(TradingCommand::CancelAllOrders(command));
1079 }
1080
1081 for order in algo_orders {
1082 self.cancel_order(order.client_order_id(), client_id, None)?;
1083 }
1084
1085 Ok(())
1086 }
1087
1088 #[expect(clippy::too_many_arguments)]
1094 fn close_position(
1095 &mut self,
1096 position: &Position,
1097 client_id: Option<ClientId>,
1098 tags: Option<Vec<Ustr>>,
1099 time_in_force: Option<TimeInForce>,
1100 reduce_only: Option<bool>,
1101 quote_quantity: Option<bool>,
1102 params: Option<Params>,
1103 ) -> anyhow::Result<()>
1104 where
1105 Self: StrategyNative,
1106 {
1107 let core = StrategyNative::strategy_core_mut(self);
1108
1109 if position.is_closed() {
1110 log::warn!("Cannot close position (already closed): {}", position.id);
1111 return Ok(());
1112 }
1113
1114 let closing_side = OrderCore::closing_side(position.side);
1115
1116 let order = core.order_factory().market(
1117 position.instrument_id,
1118 closing_side,
1119 position.quantity,
1120 time_in_force,
1121 reduce_only.or(Some(true)),
1122 quote_quantity,
1123 None,
1124 None,
1125 tags,
1126 None,
1127 );
1128
1129 self.submit_order(order, Some(position.id), client_id, params)
1130 }
1131
1132 #[expect(clippy::too_many_arguments)]
1138 fn close_all_positions(
1139 &mut self,
1140 instrument_id: InstrumentId,
1141 position_side: Option<PositionSide>,
1142 client_id: Option<ClientId>,
1143 tags: Option<Vec<Ustr>>,
1144 time_in_force: Option<TimeInForce>,
1145 reduce_only: Option<bool>,
1146 quote_quantity: Option<bool>,
1147 params: Option<Params>,
1148 ) -> anyhow::Result<()>
1149 where
1150 Self: StrategyNative,
1151 {
1152 let core = StrategyNative::strategy_core_mut(self);
1153 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
1154 let cache = core.cache_ref();
1155
1156 let positions_open = cache.positions_open(
1157 None,
1158 Some(&instrument_id),
1159 Some(&strategy_id),
1160 None,
1161 position_side,
1162 );
1163
1164 let side_str = position_side.map(|s| format!(" {s}")).unwrap_or_default();
1165
1166 if positions_open.is_empty() {
1167 log::info!("No {instrument_id} open{side_str} positions to close");
1168 return Ok(());
1169 }
1170
1171 let count = positions_open.len();
1172 log::info!(
1173 "Closing {count} open{side_str} position{}",
1174 if count == 1 { "" } else { "s" }
1175 );
1176
1177 let positions_data: Vec<_> = positions_open
1178 .iter()
1179 .map(|p| (p.id, p.instrument_id, p.side, p.quantity, p.is_closed()))
1180 .collect();
1181 drop(positions_open);
1182
1183 drop(cache);
1184
1185 for (pos_id, pos_instrument_id, pos_side, pos_quantity, is_closed) in positions_data {
1186 if is_closed {
1187 continue;
1188 }
1189
1190 let core = StrategyNative::strategy_core_mut(self);
1191 let closing_side = OrderCore::closing_side(pos_side);
1192 let order = core.order_factory().market(
1193 pos_instrument_id,
1194 closing_side,
1195 pos_quantity,
1196 time_in_force,
1197 reduce_only.or(Some(true)),
1198 quote_quantity,
1199 None,
1200 None,
1201 tags.clone(),
1202 None,
1203 );
1204
1205 self.submit_order(order, Some(pos_id), client_id, params.clone())?;
1206 }
1207
1208 Ok(())
1209 }
1210
1211 fn query_account(
1220 &mut self,
1221 account_id: AccountId,
1222 client_id: Option<ClientId>,
1223 params: Option<Params>,
1224 ) -> anyhow::Result<()>
1225 where
1226 Self: StrategyNative,
1227 {
1228 let core = StrategyNative::strategy_core_mut(self);
1229
1230 let trader_id = registered_trader_id(core)?;
1231 let ts_init = core.clock_mut().timestamp_ns();
1232
1233 let command = QueryAccount::new(
1234 trader_id,
1235 client_id,
1236 account_id,
1237 UUID4::new(),
1238 ts_init,
1239 params,
1240 None, );
1242
1243 send_exec_command(TradingCommand::QueryAccount(command));
1244 Ok(())
1245 }
1246
1247 fn query_order(
1256 &mut self,
1257 order: &OrderAny,
1258 client_id: Option<ClientId>,
1259 params: Option<Params>,
1260 ) -> anyhow::Result<()>
1261 where
1262 Self: StrategyNative,
1263 {
1264 let core = StrategyNative::strategy_core_mut(self);
1265
1266 let trader_id = registered_trader_id(core)?;
1267 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
1268 let ts_init = core.clock_mut().timestamp_ns();
1269
1270 let command = QueryOrder::new(
1271 trader_id,
1272 client_id,
1273 strategy_id,
1274 order.instrument_id(),
1275 order.client_order_id(),
1276 order.venue_order_id(),
1277 UUID4::new(),
1278 ts_init,
1279 params,
1280 None, );
1282
1283 send_exec_command(TradingCommand::QueryOrder(command));
1284 Ok(())
1285 }
1286
1287 fn handle_order_event(&mut self, event: OrderEventAny)
1289 where
1290 Self: StrategyNative,
1291 {
1292 let state = {
1293 let core = StrategyNative::strategy_core_mut(self);
1294 let id = &core.actor.actor_id;
1295 let is_warning = matches!(
1296 &event,
1297 OrderEventAny::Denied(_)
1298 | OrderEventAny::Rejected(_)
1299 | OrderEventAny::CancelRejected(_)
1300 | OrderEventAny::ModifyRejected(_)
1301 );
1302
1303 if is_warning {
1304 log::warn!("{id} {RECV}{EVT} {event}");
1305 } else if core.actor.config.log_events {
1306 log::info!("{id} {RECV}{EVT} {event}");
1307 }
1308
1309 core.actor.state()
1310 };
1311
1312 let client_order_id = event.client_order_id();
1313 let cached_order_is_closed = {
1314 let core = StrategyNative::strategy_core_mut(self);
1315 core.cache_ref()
1316 .order(&client_order_id)
1317 .map(|order| order.is_closed())
1318 };
1319 let is_terminal = match &event {
1320 OrderEventAny::FillVoided(event) => {
1321 cached_order_is_closed.unwrap_or(!event.is_reopened)
1322 }
1323 OrderEventAny::Filled(_) => cached_order_is_closed.unwrap_or(true),
1324 OrderEventAny::Canceled(_)
1325 | OrderEventAny::Rejected(_)
1326 | OrderEventAny::Expired(_)
1327 | OrderEventAny::Denied(_) => true,
1328 _ => false,
1329 };
1330
1331 if is_terminal {
1334 self.cancel_gtd_expiry(&client_order_id);
1335 }
1336
1337 if state != ComponentState::Running {
1340 return;
1341 }
1342
1343 if matches!(&event, OrderEventAny::FillVoided(event) if event.is_reopened) {
1344 let order = StrategyNative::strategy_core_mut(self)
1345 .cache_ref()
1346 .order(&client_order_id)
1347 .map(|order| order.clone());
1348 if let Some(order) = order
1349 && order.is_open()
1350 && !self.has_gtd_expiry_timer(&client_order_id)
1351 && let Err(e) = self.set_gtd_expiry(&order)
1352 {
1353 log::error!(
1354 "Failed to restore GTD expiry for reopened order {client_order_id}: {e}"
1355 );
1356 }
1357 }
1358
1359 {
1362 let core = StrategyNative::strategy_core_mut(self);
1363 if let Some(manager) = &mut core.order_manager {
1364 let actions = manager.handle_event(&event);
1365 debug_assert!(
1366 actions.is_empty(),
1367 "inactive strategy order manager returned actions"
1368 );
1369 }
1370 }
1371
1372 match &event {
1373 OrderEventAny::Initialized(e) => self.on_order_initialized(e.clone()),
1374 OrderEventAny::Denied(e) => self.on_order_denied(*e),
1375 OrderEventAny::Emulated(e) => self.on_order_emulated(*e),
1376 OrderEventAny::Released(e) => self.on_order_released(*e),
1377 OrderEventAny::Submitted(e) => self.on_order_submitted(*e),
1378 OrderEventAny::Rejected(e) => self.on_order_rejected(*e),
1379 OrderEventAny::Accepted(e) => self.on_order_accepted(*e),
1380 OrderEventAny::Canceled(e) => self.on_order_canceled(e),
1381 OrderEventAny::Expired(e) => self.on_order_expired(*e),
1382 OrderEventAny::Triggered(e) => self.on_order_triggered(*e),
1383 OrderEventAny::PendingUpdate(e) => self.on_order_pending_update(*e),
1384 OrderEventAny::PendingCancel(e) => self.on_order_pending_cancel(*e),
1385 OrderEventAny::ModifyRejected(e) => self.on_order_modify_rejected(*e),
1386 OrderEventAny::CancelRejected(e) => self.on_order_cancel_rejected(*e),
1387 OrderEventAny::Updated(e) => self.on_order_updated(*e),
1388 OrderEventAny::Filled(e) => self.on_order_filled(e),
1389 OrderEventAny::FillVoided(e) => self.on_order_fill_voided(e),
1390 }
1391 self.on_order_event(event);
1392 }
1393
1394 fn handle_position_event(&mut self, event: PositionEvent)
1396 where
1397 Self: StrategyNative,
1398 {
1399 let state = {
1400 let core = StrategyNative::strategy_core_mut(self);
1401
1402 if core.actor.config.log_events {
1403 let id = &core.actor.actor_id;
1404 log::info!("{id} {RECV}{EVT} {event:?}");
1405 }
1406
1407 core.actor.state()
1408 };
1409
1410 if state != ComponentState::Running {
1411 return;
1412 }
1413
1414 match &event {
1415 PositionEvent::PositionOpened(e) => self.on_position_opened(e.clone()),
1416 PositionEvent::PositionChanged(e) => self.on_position_changed(e.clone()),
1417 PositionEvent::PositionClosed(e) => self.on_position_closed(e.clone()),
1418 PositionEvent::PositionAdjusted(_) => {
1419 return;
1420 }
1421 }
1422 self.on_position_event(event);
1423 }
1424
1425 fn on_start(&mut self) -> anyhow::Result<()>
1436 where
1437 Self: StrategyNative,
1438 {
1439 let core = StrategyNative::strategy_core_mut(self);
1440 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
1441 log::info!("Starting {strategy_id}");
1442
1443 if core.config.manage_gtd_expiry {
1444 self.reactivate_gtd_timers();
1445 }
1446
1447 Ok(())
1448 }
1449
1450 fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()>
1459 where
1460 Self: StrategyNative,
1461 {
1462 if event.name.starts_with("GTD-EXPIRY:") {
1463 self.expire_gtd_order(event.clone());
1464 } else if event.name.starts_with("MARKET_EXIT_CHECK:") {
1465 self.check_market_exit(event.clone());
1466 }
1467 Ok(())
1468 }
1469
1470 #[allow(unused_variables)]
1476 fn on_order_initialized(&mut self, event: OrderInitialized) {}
1477
1478 #[allow(unused_variables)]
1482 fn on_order_event(&mut self, event: OrderEventAny) {}
1483
1484 #[allow(unused_variables)]
1488 fn on_order_denied(&mut self, event: OrderDenied) {}
1489
1490 #[allow(unused_variables)]
1494 fn on_order_emulated(&mut self, event: OrderEmulated) {}
1495
1496 #[allow(unused_variables)]
1500 fn on_order_released(&mut self, event: OrderReleased) {}
1501
1502 #[allow(unused_variables)]
1506 fn on_order_submitted(&mut self, event: OrderSubmitted) {}
1507
1508 #[allow(unused_variables)]
1512 fn on_order_rejected(&mut self, event: OrderRejected) {}
1513
1514 #[allow(unused_variables)]
1518 fn on_order_accepted(&mut self, event: OrderAccepted) {}
1519
1520 #[allow(unused_variables)]
1524 fn on_order_expired(&mut self, event: OrderExpired) {}
1525
1526 #[allow(unused_variables)]
1530 fn on_order_triggered(&mut self, event: OrderTriggered) {}
1531
1532 #[allow(unused_variables)]
1536 fn on_order_pending_update(&mut self, event: OrderPendingUpdate) {}
1537
1538 #[allow(unused_variables)]
1542 fn on_order_pending_cancel(&mut self, event: OrderPendingCancel) {}
1543
1544 #[allow(unused_variables)]
1548 fn on_order_modify_rejected(&mut self, event: OrderModifyRejected) {}
1549
1550 #[allow(unused_variables)]
1554 fn on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {}
1555
1556 #[allow(unused_variables)]
1560 fn on_order_updated(&mut self, event: OrderUpdated) {}
1561
1562 #[allow(unused_variables)]
1566 fn on_order_canceled(&mut self, event: &OrderCanceled) {}
1567
1568 #[allow(unused_variables)]
1572 fn on_order_filled(&mut self, event: &OrderFilled) {}
1573
1574 #[allow(unused_variables)]
1576 fn on_order_fill_voided(&mut self, event: &OrderFillVoided) {}
1577
1578 #[allow(unused_variables)]
1582 fn on_position_opened(&mut self, event: PositionOpened) {}
1583
1584 #[allow(unused_variables)]
1588 fn on_position_event(&mut self, event: PositionEvent) {}
1589
1590 #[allow(unused_variables)]
1594 fn on_position_changed(&mut self, event: PositionChanged) {}
1595
1596 #[allow(unused_variables)]
1600 fn on_position_closed(&mut self, event: PositionClosed) {}
1601
1602 fn on_market_exit(&mut self) {}
1606
1607 fn post_market_exit(&mut self) {}
1611
1612 fn is_exiting(&self) -> bool
1616 where
1617 Self: StrategyNative,
1618 {
1619 StrategyNative::strategy_core(self).is_exiting
1620 }
1621
1622 fn market_exit(&mut self) -> anyhow::Result<()>
1638 where
1639 Self: StrategyNative,
1640 {
1641 let core = StrategyNative::strategy_core_mut(self);
1642 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
1643
1644 if core.actor.state() != ComponentState::Running {
1645 log::warn!("{strategy_id} Cannot market exit: strategy is not running");
1646 return Ok(());
1647 }
1648
1649 if core.is_exiting {
1650 log::warn!("{strategy_id} Market exit called when already in progress");
1651 return Ok(());
1652 }
1653
1654 core.is_exiting = true;
1655 core.market_exit_attempts = 0;
1656 let time_in_force = core.config.market_exit_time_in_force;
1657 let reduce_only = core.config.market_exit_reduce_only;
1658
1659 log::info!("{strategy_id} Initiating market exit...");
1660
1661 self.on_market_exit();
1662
1663 let core = StrategyNative::strategy_core_mut(self);
1664 let cache = core.cache_ref();
1665
1666 let mut instruments: AHashSet<InstrumentId> = AHashSet::new();
1667
1668 for client_order_id in
1669 cache.iter_client_order_ids_open(None, None, Some(&strategy_id), None)
1670 {
1671 if let Some(order) = cache.order(&client_order_id) {
1672 instruments.insert(order.instrument_id());
1673 }
1674 }
1675
1676 for client_order_id in
1677 cache.iter_client_order_ids_inflight(None, None, Some(&strategy_id), None)
1678 {
1679 if let Some(order) = cache.order(&client_order_id) {
1680 instruments.insert(order.instrument_id());
1681 }
1682 }
1683
1684 for position_id in cache.iter_position_open_ids(None, None, Some(&strategy_id), None) {
1685 if let Some(position) = cache.position(&position_id) {
1686 instruments.insert(position.instrument_id);
1687 }
1688 }
1689
1690 let market_exit_tag = core.market_exit_tag;
1691 let mut instruments: Vec<_> = instruments.into_iter().collect();
1695 instruments.sort();
1696 drop(cache);
1697
1698 for instrument_id in instruments {
1699 if let Err(e) = self.cancel_all_orders(instrument_id, None, None, None) {
1700 log::error!("Error canceling orders for {instrument_id}: {e}");
1701 }
1702
1703 if let Err(e) = self.close_all_positions(
1704 instrument_id,
1705 None,
1706 None,
1707 Some(vec![market_exit_tag]),
1708 Some(time_in_force),
1709 Some(reduce_only),
1710 None,
1711 None,
1712 ) {
1713 log::error!("Error closing positions for {instrument_id}: {e}");
1714 }
1715 }
1716
1717 let core = StrategyNative::strategy_core_mut(self);
1718 let interval_ms = core.config.market_exit_interval_ms;
1719 let timer_name = core.market_exit_timer_name;
1720
1721 log::info!("{strategy_id} Setting market exit timer at {interval_ms}ms intervals");
1722
1723 let interval_ns = interval_ms * 1_000_000;
1724 let result = core.clock_mut().set_timer_ns(
1725 timer_name.as_str(),
1726 interval_ns,
1727 None,
1728 None,
1729 None,
1730 None,
1731 None,
1732 );
1733
1734 if let Err(e) = result {
1735 core.is_exiting = false;
1737 core.market_exit_attempts = 0;
1738 return Err(e);
1739 }
1740
1741 Ok(())
1742 }
1743
1744 fn check_market_exit(&mut self, _event: TimeEvent)
1748 where
1749 Self: StrategyNative,
1750 {
1751 if !self.is_exiting() {
1753 return;
1754 }
1755
1756 let core = StrategyNative::strategy_core_mut(self);
1757 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
1758
1759 core.market_exit_attempts += 1;
1760 let attempts = core.market_exit_attempts;
1761 let max_attempts = core.config.market_exit_max_attempts;
1762
1763 log::debug!(
1764 "{strategy_id} Market exit check triggered (attempt {attempts}/{max_attempts})"
1765 );
1766
1767 if attempts >= max_attempts {
1768 let cache = core.cache_ref();
1769 let open_orders_count =
1770 cache.orders_open_count(None, None, Some(&strategy_id), None, None);
1771 let inflight_orders_count =
1772 cache.orders_inflight_count(None, None, Some(&strategy_id), None, None);
1773 let open_positions_count =
1774 cache.positions_open_count(None, None, Some(&strategy_id), None, None);
1775
1776 drop(cache);
1777
1778 log::warn!(
1779 "{strategy_id} Market exit max attempts ({max_attempts}) reached, \
1780 completing with open orders: {open_orders_count}, \
1781 inflight orders: {inflight_orders_count}, \
1782 open positions: {open_positions_count}"
1783 );
1784
1785 self.finalize_market_exit();
1786 return;
1787 }
1788
1789 let cache = core.cache_ref();
1790 let has_open_orders = !cache
1791 .orders_open(None, None, Some(&strategy_id), None, None)
1792 .is_empty();
1793 let has_inflight_orders = !cache
1794 .orders_inflight(None, None, Some(&strategy_id), None, None)
1795 .is_empty();
1796
1797 if has_open_orders || has_inflight_orders {
1798 return;
1799 }
1800
1801 let positions_data: Vec<_> = cache
1802 .positions_open(None, None, Some(&strategy_id), None, None)
1803 .iter()
1804 .map(|p| (p.id, p.instrument_id, p.side, p.quantity, p.is_closed()))
1805 .collect();
1806
1807 if !positions_data.is_empty() {
1808 drop(cache);
1810
1811 for (pos_id, instrument_id, side, quantity, is_closed) in positions_data {
1812 if is_closed {
1813 continue;
1814 }
1815
1816 let core = StrategyNative::strategy_core_mut(self);
1817 let time_in_force = core.config.market_exit_time_in_force;
1818 let reduce_only = core.config.market_exit_reduce_only;
1819 let market_exit_tag = core.market_exit_tag;
1820 let closing_side = OrderCore::closing_side(side);
1821 let order = core.order_factory().market(
1822 instrument_id,
1823 closing_side,
1824 quantity,
1825 Some(time_in_force),
1826 Some(reduce_only),
1827 None,
1828 None,
1829 None,
1830 Some(vec![market_exit_tag]),
1831 None,
1832 );
1833
1834 if let Err(e) = self.submit_order(order, Some(pos_id), None, None) {
1835 log::error!("Error re-submitting close order for position {pos_id}: {e}");
1836 }
1837 }
1838 return;
1839 }
1840
1841 drop(cache);
1842 self.finalize_market_exit();
1843 }
1844
1845 fn finalize_market_exit(&mut self)
1850 where
1851 Self: StrategyNative,
1852 {
1853 let (strategy_id, should_stop) = {
1854 let core = StrategyNative::strategy_core_mut(self);
1855 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
1856 let should_stop = core.pending_stop;
1857 (strategy_id, should_stop)
1858 };
1859
1860 self.cancel_market_exit();
1861
1862 let hook_result = catch_unwind(AssertUnwindSafe(|| {
1863 self.post_market_exit();
1864 }));
1865
1866 if let Err(e) = hook_result {
1867 log::error!("{strategy_id} Error in post_market_exit: {e:?}");
1868 }
1869
1870 if should_stop {
1871 log::info!("{strategy_id} Market exit complete, stopping strategy");
1872
1873 if let Err(e) = Component::stop(self) {
1874 log::error!("{strategy_id} Failed to stop: {e}");
1875 }
1876 }
1877
1878 let core = StrategyNative::strategy_core_mut(self);
1879 debug_assert!(
1880 !(core.pending_stop
1881 && !core.is_exiting
1882 && core.actor.state() == ComponentState::Running),
1883 "INVARIANT: stuck state after finalize_market_exit"
1884 );
1885 }
1886
1887 fn cancel_market_exit(&mut self)
1891 where
1892 Self: StrategyNative,
1893 {
1894 let core = StrategyNative::strategy_core_mut(self);
1895 let timer_name = core.market_exit_timer_name;
1896
1897 if core
1898 .clock_mut()
1899 .timer_names()
1900 .contains(&timer_name.as_str())
1901 {
1902 core.clock_mut().cancel_timer(timer_name.as_str());
1903 }
1904
1905 core.is_exiting = false;
1906 core.pending_stop = false;
1907 core.market_exit_attempts = 0;
1908 }
1909
1910 fn stop(&mut self) -> bool
1922 where
1923 Self: StrategyNative,
1924 {
1925 let (manage_stop, is_exiting, should_initiate_exit) = {
1926 let core = StrategyNative::strategy_core_mut(self);
1927 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
1928 let manage_stop = core.config.manage_stop;
1929 let state = core.actor.state();
1930 let pending_stop = core.pending_stop;
1931 let is_exiting = core.is_exiting;
1932
1933 if manage_stop {
1934 if state != ComponentState::Running {
1935 return true; }
1937
1938 if pending_stop {
1939 return false; }
1941
1942 core.pending_stop = true;
1943 let should_initiate_exit = !is_exiting;
1944
1945 if should_initiate_exit {
1946 log::info!("{strategy_id} Initiating market exit before stop");
1947 }
1948
1949 (manage_stop, is_exiting, should_initiate_exit)
1950 } else {
1951 (manage_stop, is_exiting, false)
1952 }
1953 };
1954
1955 if manage_stop {
1956 if should_initiate_exit && let Err(e) = self.market_exit() {
1957 log::warn!("Market exit failed during stop: {e}, proceeding with stop");
1958 StrategyNative::strategy_core_mut(self).pending_stop = false;
1959 return true;
1960 }
1961 debug_assert!(
1962 self.is_exiting(),
1963 "INVARIANT: deferring stop but not exiting"
1964 );
1965 return false; }
1967
1968 if is_exiting {
1970 self.cancel_market_exit();
1971 }
1972
1973 true }
1975
1976 fn deny_order(&mut self, order: &OrderAny, reason: Ustr)
1981 where
1982 Self: StrategyNative,
1983 {
1984 let core = StrategyNative::strategy_core_mut(self);
1985 let Some(trader_id) = core.trader_id() else {
1986 log::error!(
1987 "Cannot deny order {}: trader_id is not set",
1988 order.client_order_id()
1989 );
1990 return;
1991 };
1992 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
1993 let ts_now = core.clock_mut().timestamp_ns();
1994
1995 let event = OrderDenied::new(
1996 trader_id,
1997 strategy_id,
1998 order.instrument_id(),
1999 order.client_order_id(),
2000 reason,
2001 UUID4::new(),
2002 ts_now,
2003 ts_now,
2004 );
2005
2006 log::warn!(
2007 "{strategy_id} Order {} denied: {reason}",
2008 order.client_order_id()
2009 );
2010
2011 let publish_initialized = {
2012 let cache_rc = core.cache_rc();
2013 let mut cache = cache_rc.borrow_mut();
2014 if cache.order_exists(&order.client_order_id()) {
2015 false
2016 } else {
2017 match cache.add_order(order.clone(), None, None, true) {
2018 Ok(()) => true,
2019 Err(e) => {
2020 log::warn!("Failed to add denied order to cache: {e}");
2021 false
2022 }
2023 }
2024 }
2025 };
2026
2027 if publish_initialized {
2028 publish_order_initialized(order);
2029 }
2030
2031 let event = OrderEventAny::Denied(event);
2032 let applied = {
2033 let cache_rc = core.cache_rc();
2034 let mut cache = cache_rc.borrow_mut();
2035 if let Err(e) = cache.update_order(&event) {
2036 log::warn!("Failed to apply OrderDenied event: {e}");
2037 false
2038 } else {
2039 true
2040 }
2041 };
2042
2043 if applied {
2044 let topic = format!("events.order.{strategy_id}");
2045 msgbus::publish_order_event(topic.into(), &event);
2046 }
2047 }
2048
2049 fn deny_order_list(&mut self, orders: &[OrderAny], reason: Ustr)
2053 where
2054 Self: StrategyNative,
2055 {
2056 for order in orders {
2057 if !order.is_closed() {
2058 self.deny_order(order, reason);
2059 }
2060 }
2061 }
2062
2063 fn set_gtd_expiry(&mut self, order: &OrderAny) -> anyhow::Result<()>
2073 where
2074 Self: StrategyNative,
2075 {
2076 let core = StrategyNative::strategy_core_mut(self);
2077
2078 if !core.config.manage_gtd_expiry || order.time_in_force() != TimeInForce::Gtd {
2079 return Ok(());
2080 }
2081
2082 let Some(expire_time) = order.expire_time() else {
2083 return Ok(());
2084 };
2085
2086 let client_order_id = order.client_order_id();
2087 let timer_name = format!("GTD-EXPIRY:{client_order_id}");
2088
2089 let current_time_ns = {
2090 let clock = core.clock_mut();
2091 clock.timestamp_ns()
2092 };
2093
2094 if current_time_ns >= expire_time.as_u64() {
2095 log::info!("GTD order {client_order_id} already expired, canceling immediately");
2096 return self.cancel_order(order.client_order_id(), None, None);
2097 }
2098
2099 {
2100 let mut clock = core.clock_mut();
2101 clock.set_time_alert_ns(&timer_name, expire_time, None, None)?;
2102 }
2103
2104 core.gtd_timers
2105 .insert(client_order_id, Ustr::from(&timer_name));
2106
2107 log::debug!("Set GTD expiry timer for {client_order_id} at {expire_time}");
2108 Ok(())
2109 }
2110
2111 fn cancel_gtd_expiry(&mut self, client_order_id: &ClientOrderId)
2113 where
2114 Self: StrategyNative,
2115 {
2116 let core = StrategyNative::strategy_core_mut(self);
2117
2118 if let Some(timer_name) = core.gtd_timers.remove(client_order_id) {
2119 core.clock_mut().cancel_timer(timer_name.as_str());
2120 log::debug!("Canceled GTD expiry timer for {client_order_id}");
2121 }
2122 }
2123
2124 fn has_gtd_expiry_timer(&mut self, client_order_id: &ClientOrderId) -> bool
2126 where
2127 Self: StrategyNative,
2128 {
2129 let core = StrategyNative::strategy_core_mut(self);
2130 core.gtd_timers.contains_key(client_order_id)
2131 }
2132
2133 fn expire_gtd_order(&mut self, event: TimeEvent)
2137 where
2138 Self: StrategyNative,
2139 {
2140 let timer_name = event.name.to_string();
2141 let Some(client_order_id_str) = timer_name.strip_prefix("GTD-EXPIRY:") else {
2142 log::error!("Invalid GTD timer name format: {timer_name}");
2143 return;
2144 };
2145
2146 let client_order_id = ClientOrderId::from(client_order_id_str);
2147
2148 let core = StrategyNative::strategy_core_mut(self);
2149 core.gtd_timers.remove(&client_order_id);
2150
2151 let order = core.cache_ref().order(&client_order_id).map(|o| o.clone());
2152 let Some(order) = order else {
2153 log::warn!("GTD order {client_order_id} not found in cache");
2154 return;
2155 };
2156
2157 log::info!("GTD order {client_order_id} expired");
2158
2159 if let Err(e) = self.cancel_order(order.client_order_id(), None, None) {
2160 log::error!("Failed to cancel expired GTD order {client_order_id}: {e}");
2161 }
2162 }
2163
2164 fn reactivate_gtd_timers(&mut self)
2169 where
2170 Self: StrategyNative,
2171 {
2172 let core = StrategyNative::strategy_core_mut(self);
2173 let strategy_id = StrategyId::from(core.actor_id().inner().as_str());
2174 let current_time_ns = core.clock_mut().timestamp_ns();
2175
2176 let gtd_orders: Vec<OrderAny> = core
2177 .cache_ref()
2178 .orders_open(None, None, Some(&strategy_id), None, None)
2179 .into_iter()
2180 .filter(|o| o.time_in_force() == TimeInForce::Gtd)
2181 .map(|o| o.clone())
2182 .collect();
2183
2184 for order in gtd_orders {
2185 let Some(expire_time) = order.expire_time() else {
2186 continue;
2187 };
2188
2189 let expire_time_ns = expire_time.as_u64();
2190 let client_order_id = order.client_order_id();
2191
2192 if current_time_ns >= expire_time_ns {
2193 log::info!("GTD order {client_order_id} already expired, canceling immediately");
2194 if let Err(e) = self.cancel_order(order.client_order_id(), None, None) {
2195 log::error!("Failed to cancel expired GTD order {client_order_id}: {e}");
2196 }
2197 } else if let Err(e) = self.set_gtd_expiry(&order) {
2198 log::error!("Failed to set GTD expiry timer for {client_order_id}: {e}");
2199 }
2200 }
2201 }
2202}
2203
2204fn publish_order_initialized(order: &OrderAny) {
2205 let topic = format!("events.order.{}", order.strategy_id());
2206 let event = OrderEventAny::Initialized(order.init_event().clone());
2207 msgbus::publish_order_event(topic.into(), &event);
2208}
2209
2210fn send_emulator_command(command: TradingCommand) {
2211 log_cmd_send(&command);
2212 let endpoint = MessagingSwitchboard::order_emulator_execute();
2213 msgbus::send_trading_command(endpoint, command);
2214}
2215
2216fn send_algo_command(command: SubmitOrder, exec_algorithm_id: ExecAlgorithmId) {
2217 let id = command.strategy_id;
2218 log::info!("{id} {CMD}{SEND} {command}");
2219
2220 let endpoint = format!("{exec_algorithm_id}.execute");
2221 msgbus::send_any(endpoint.into(), &TradingCommand::SubmitOrder(command));
2222}
2223
2224fn send_risk_command(command: TradingCommand) {
2225 log_cmd_send(&command);
2226 let endpoint = MessagingSwitchboard::risk_engine_queue_execute();
2227 msgbus::send_trading_command(endpoint, command);
2228}
2229
2230fn send_exec_command(command: TradingCommand) {
2231 log_cmd_send(&command);
2232 let endpoint = MessagingSwitchboard::exec_engine_queue_execute();
2233 msgbus::send_trading_command(endpoint, command);
2234}
2235
2236fn log_cmd_send(command: &TradingCommand) {
2237 if let Some(id) = command.strategy_id() {
2238 log::info!("{id} {CMD}{SEND} {command}");
2239 } else {
2240 log::info!("{CMD}{SEND} {command}");
2241 }
2242}
2243
2244fn registered_trader_id(core: &StrategyCore) -> anyhow::Result<TraderId> {
2245 core.trader_id()
2246 .ok_or_else(|| anyhow::anyhow!("Strategy not registered: trader_id is not set"))
2247}
2248
2249fn required_account_id(order: &OrderAny, operation: &str) -> anyhow::Result<AccountId> {
2250 order.account_id().ok_or_else(|| {
2251 anyhow::anyhow!(
2252 "Cannot generate {operation} event for {}: account_id is not set",
2253 order.client_order_id()
2254 )
2255 })
2256}
2257
2258#[cfg(test)]
2259mod tests {
2260 use std::{cell::RefCell, rc::Rc};
2261
2262 use nautilus_common::{
2263 actor::DataActor,
2264 cache::{Cache, ORDER_NOT_FOUND},
2265 clock::{Clock, TestClock},
2266 component::Component,
2267 enums::{ComponentState, ComponentTrigger},
2268 msgbus::{
2269 self, MessagingSwitchboard, TypedHandler, TypedIntoHandler,
2270 stubs::{
2271 TypedIntoMessageSavingHandler, TypedMessageSavingHandler,
2272 get_typed_into_message_saving_handler, get_typed_message_saving_handler,
2273 },
2274 },
2275 timer::{TimeEvent, TimeEventCallback},
2276 };
2277 use nautilus_core::UnixNanos;
2278 use nautilus_model::{
2279 enums::{
2280 LiquiditySide, OrderSide, OrderStatus, OrderType, PositionAdjustmentType, PositionSide,
2281 },
2282 events::{
2283 OrderAccepted, OrderCanceled, OrderFilled, OrderRejected, PositionAdjusted,
2284 order::spec::{
2285 OrderAcceptedSpec, OrderCanceledSpec, OrderExpiredSpec, OrderFillVoidedSpec,
2286 OrderFilledSpec, OrderRejectedSpec,
2287 },
2288 },
2289 identifiers::{
2290 AccountId, ActorId, ClientOrderId, ComponentId, InstrumentId, OrderListId, PositionId,
2291 StrategyId, TradeId, TraderId, VenueOrderId,
2292 },
2293 orderbook::own::OwnOrderBook,
2294 orders::{LimitOrder, MarketOrder, OrderTestBuilder, stubs::TestOrderEventStubs},
2295 stubs::TestDefault,
2296 types::{Currency, Money, Price},
2297 };
2298 use nautilus_portfolio::portfolio::Portfolio;
2299 use rstest::rstest;
2300 use serde_json::Value;
2301
2302 use super::*;
2303 use crate::nautilus_strategy;
2304
2305 #[derive(Debug)]
2306 struct TestStrategy {
2307 core: StrategyCore,
2308 on_order_rejected_called: bool,
2309 on_order_event_called: bool,
2310 on_order_accepted_called: bool,
2311 on_order_canceled_called: bool,
2312 on_order_filled_called: bool,
2313 on_order_fill_voided_called: bool,
2314 on_order_expired_called: bool,
2315 on_position_event_called: bool,
2316 on_position_opened_called: bool,
2317 on_position_changed_called: bool,
2318 on_position_closed_called: bool,
2319 }
2320
2321 #[derive(Debug)]
2322 struct CoreFreeStrategy {
2323 state: ComponentState,
2324 started: bool,
2325 }
2326
2327 impl Component for CoreFreeStrategy {
2328 fn component_id(&self) -> ComponentId {
2329 ComponentId::new("CoreFreeStrategy")
2330 }
2331
2332 fn state(&self) -> ComponentState {
2333 self.state
2334 }
2335
2336 fn transition_state(&mut self, trigger: ComponentTrigger) -> anyhow::Result<()> {
2337 self.state = self.state.transition(&trigger)?;
2338 Ok(())
2339 }
2340
2341 fn register(
2342 &mut self,
2343 _trader_id: TraderId,
2344 _clock: Rc<RefCell<dyn Clock>>,
2345 _cache: Rc<RefCell<Cache>>,
2346 ) -> anyhow::Result<()> {
2347 Ok(())
2348 }
2349 }
2350
2351 impl DataActor for CoreFreeStrategy {
2352 fn on_start(&mut self) -> anyhow::Result<()> {
2353 self.started = true;
2354 Ok(())
2355 }
2356 }
2357
2358 impl Strategy for CoreFreeStrategy {}
2359
2360 impl TestStrategy {
2361 fn new(config: StrategyConfig) -> Self {
2362 Self {
2363 core: StrategyCore::new(config),
2364 on_order_rejected_called: false,
2365 on_order_event_called: false,
2366 on_order_accepted_called: false,
2367 on_order_canceled_called: false,
2368 on_order_filled_called: false,
2369 on_order_fill_voided_called: false,
2370 on_order_expired_called: false,
2371 on_position_event_called: false,
2372 on_position_opened_called: false,
2373 on_position_changed_called: false,
2374 on_position_closed_called: false,
2375 }
2376 }
2377 }
2378
2379 impl DataActor for TestStrategy {}
2380
2381 nautilus_strategy!(TestStrategy, {
2382 fn on_order_canceled(&mut self, _event: &OrderCanceled) {
2383 self.on_order_canceled_called = true;
2384 }
2385
2386 fn on_order_filled(&mut self, _event: &OrderFilled) {
2387 self.on_order_filled_called = true;
2388 }
2389
2390 fn on_order_fill_voided(&mut self, _event: &OrderFillVoided) {
2391 self.on_order_fill_voided_called = true;
2392 }
2393
2394 fn on_order_rejected(&mut self, _event: OrderRejected) {
2395 self.on_order_rejected_called = true;
2396 }
2397
2398 fn on_order_event(&mut self, _event: OrderEventAny) {
2399 self.on_order_event_called = true;
2400 }
2401
2402 fn on_order_accepted(&mut self, _event: OrderAccepted) {
2403 self.on_order_accepted_called = true;
2404 }
2405
2406 fn on_order_expired(&mut self, _event: OrderExpired) {
2407 self.on_order_expired_called = true;
2408 }
2409
2410 fn on_position_opened(&mut self, _event: PositionOpened) {
2411 self.on_position_opened_called = true;
2412 }
2413
2414 fn on_position_event(&mut self, _event: PositionEvent) {
2415 self.on_position_event_called = true;
2416 }
2417
2418 fn on_position_changed(&mut self, _event: PositionChanged) {
2419 self.on_position_changed_called = true;
2420 }
2421
2422 fn on_position_closed(&mut self, _event: PositionClosed) {
2423 self.on_position_closed_called = true;
2424 }
2425 });
2426
2427 fn create_test_strategy() -> TestStrategy {
2428 let config = StrategyConfig {
2429 strategy_id: Some(StrategyId::from("TEST-001")),
2430 order_id_tag: Some("001".to_string()),
2431 ..Default::default()
2432 };
2433 TestStrategy::new(config)
2434 }
2435
2436 fn register_strategy(strategy: &mut TestStrategy) {
2437 let trader_id = TraderId::from("TRADER-001");
2438 let clock = Rc::new(RefCell::new(TestClock::new()));
2439 let cache = Rc::new(RefCell::new(Cache::default()));
2440 let portfolio = Rc::new(RefCell::new(Portfolio::new(
2441 clock.clone(),
2442 cache.clone(),
2443 None,
2444 )));
2445
2446 strategy
2447 .core
2448 .register(trader_id, clock, cache, portfolio)
2449 .unwrap();
2450 strategy.initialize().unwrap();
2451 }
2452
2453 fn start_strategy(strategy: &mut TestStrategy) {
2454 strategy.start().unwrap();
2455 }
2456
2457 fn stop_strategy(strategy: &mut TestStrategy) {
2458 Component::stop(strategy).unwrap();
2459 }
2460
2461 fn make_filled(client_order_id: ClientOrderId) -> OrderEventAny {
2462 OrderEventAny::Filled(
2463 OrderFilledSpec::builder()
2464 .trader_id(TraderId::from("TRADER-001"))
2465 .strategy_id(StrategyId::from("TEST-001"))
2466 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2467 .client_order_id(client_order_id)
2468 .venue_order_id(VenueOrderId::test_default())
2469 .account_id(AccountId::from("ACC-001"))
2470 .trade_id(TradeId::test_default())
2471 .last_qty(Quantity::default())
2472 .last_px(Price::default())
2473 .currency(Currency::from("USD"))
2474 .liquidity_side(LiquiditySide::Taker)
2475 .event_id(UUID4::default())
2476 .build(),
2477 )
2478 }
2479
2480 fn make_fill_voided(client_order_id: ClientOrderId, is_reopened: bool) -> OrderEventAny {
2481 OrderEventAny::FillVoided(
2482 OrderFillVoidedSpec::builder()
2483 .trader_id(TraderId::from("TRADER-001"))
2484 .strategy_id(StrategyId::from("TEST-001"))
2485 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2486 .client_order_id(client_order_id)
2487 .venue_order_id(VenueOrderId::test_default())
2488 .account_id(AccountId::from("ACC-001"))
2489 .is_reopened(is_reopened)
2490 .build(),
2491 )
2492 }
2493
2494 fn make_terminal_fill_voided(client_order_id: ClientOrderId) -> OrderEventAny {
2495 make_fill_voided(client_order_id, false)
2496 }
2497
2498 fn make_canceled(client_order_id: ClientOrderId) -> OrderEventAny {
2499 OrderEventAny::Canceled(
2500 OrderCanceledSpec::builder()
2501 .trader_id(TraderId::from("TRADER-001"))
2502 .strategy_id(StrategyId::from("TEST-001"))
2503 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2504 .client_order_id(client_order_id)
2505 .account_id(AccountId::from("ACC-001"))
2506 .event_id(UUID4::default())
2507 .build(),
2508 )
2509 }
2510
2511 fn make_rejected(client_order_id: ClientOrderId) -> OrderEventAny {
2512 OrderEventAny::Rejected(
2513 OrderRejectedSpec::builder()
2514 .trader_id(TraderId::from("TRADER-001"))
2515 .strategy_id(StrategyId::from("TEST-001"))
2516 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2517 .client_order_id(client_order_id)
2518 .account_id(AccountId::from("ACC-001"))
2519 .reason("Test rejection".into())
2520 .event_id(UUID4::default())
2521 .build(),
2522 )
2523 }
2524
2525 fn make_expired(client_order_id: ClientOrderId) -> OrderEventAny {
2526 OrderEventAny::Expired(
2527 OrderExpiredSpec::builder()
2528 .trader_id(TraderId::from("TRADER-001"))
2529 .strategy_id(StrategyId::from("TEST-001"))
2530 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2531 .client_order_id(client_order_id)
2532 .account_id(AccountId::from("ACC-001"))
2533 .event_id(UUID4::default())
2534 .build(),
2535 )
2536 }
2537
2538 fn make_accepted(client_order_id: ClientOrderId) -> OrderEventAny {
2539 OrderEventAny::Accepted(
2540 OrderAcceptedSpec::builder()
2541 .trader_id(TraderId::from("TRADER-001"))
2542 .strategy_id(StrategyId::from("TEST-001"))
2543 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2544 .client_order_id(client_order_id)
2545 .venue_order_id(VenueOrderId::test_default())
2546 .account_id(AccountId::from("ACC-001"))
2547 .event_id(UUID4::default())
2548 .build(),
2549 )
2550 }
2551
2552 fn make_accepted_market_order(client_order_id: &str) -> OrderAny {
2553 let mut order = OrderAny::Market(MarketOrder::new(
2554 TraderId::from("TRADER-001"),
2555 StrategyId::from("TEST-001"),
2556 InstrumentId::from("BTCUSDT.BINANCE"),
2557 ClientOrderId::from(client_order_id),
2558 OrderSide::Buy,
2559 Quantity::from(100_000),
2560 TimeInForce::Gtc,
2561 UUID4::new(),
2562 UnixNanos::default(),
2563 false,
2564 false,
2565 None,
2566 None,
2567 None,
2568 None,
2569 None,
2570 None,
2571 None,
2572 None,
2573 ));
2574 let account_id = AccountId::from("ACC-001");
2575 order
2576 .apply(TestOrderEventStubs::submitted(&order, account_id))
2577 .unwrap();
2578 order
2579 .apply(TestOrderEventStubs::accepted(
2580 &order,
2581 account_id,
2582 VenueOrderId::from(client_order_id),
2585 ))
2586 .unwrap();
2587 order
2588 }
2589
2590 fn make_accepted_limit_order(client_order_id: &str) -> OrderAny {
2591 let mut order = OrderAny::Limit(LimitOrder::new(
2592 TraderId::from("TRADER-001"),
2593 StrategyId::from("TEST-001"),
2594 InstrumentId::from("BTCUSDT.BINANCE"),
2595 ClientOrderId::from(client_order_id),
2596 OrderSide::Buy,
2597 Quantity::from("1.0"),
2598 Price::from("50000.0"),
2599 TimeInForce::Gtc,
2600 None,
2601 false,
2602 false,
2603 false,
2604 None,
2605 None,
2606 None,
2607 None,
2608 None,
2609 None,
2610 None,
2611 None,
2612 None,
2613 None,
2614 None,
2615 UUID4::new(),
2616 UnixNanos::default(),
2617 ));
2618 let account_id = AccountId::from("ACC-001");
2619 order
2620 .apply(TestOrderEventStubs::submitted(&order, account_id))
2621 .unwrap();
2622 order
2623 .apply(TestOrderEventStubs::accepted(
2624 &order,
2625 account_id,
2626 VenueOrderId::from(client_order_id),
2628 ))
2629 .unwrap();
2630 order
2631 }
2632
2633 fn make_initialized_market_order(client_order_id: &str) -> OrderAny {
2634 OrderAny::Market(MarketOrder::new(
2635 TraderId::from("TRADER-001"),
2636 StrategyId::from("TEST-001"),
2637 InstrumentId::from("BTCUSDT.BINANCE"),
2638 ClientOrderId::from(client_order_id),
2639 OrderSide::Buy,
2640 Quantity::from(100_000),
2641 TimeInForce::Gtc,
2642 UUID4::new(),
2643 UnixNanos::default(),
2644 false,
2645 false,
2646 None,
2647 None,
2648 None,
2649 None,
2650 None,
2651 None,
2652 None,
2653 None,
2654 ))
2655 }
2656
2657 fn add_order_to_cache(strategy: &TestStrategy, order: &OrderAny) {
2658 let cache_rc = strategy.core.cache_rc();
2659 let mut cache = cache_rc.borrow_mut();
2660 cache.add_order(order.clone(), None, None, true).unwrap();
2661 }
2662
2663 fn add_order_to_cache_and_own_book(strategy: &TestStrategy, order: &OrderAny) {
2664 let cache_rc = strategy.core.cache_rc();
2665 let mut cache = cache_rc.borrow_mut();
2666 cache.add_order(order.clone(), None, None, true).unwrap();
2667 cache
2668 .add_own_order_book(OwnOrderBook::new(order.instrument_id()))
2669 .unwrap();
2670 cache.update_own_order_book(order);
2671 }
2672
2673 fn make_position_opened() -> PositionEvent {
2674 PositionEvent::PositionOpened(PositionOpened {
2675 trader_id: TraderId::from("TRADER-001"),
2676 strategy_id: StrategyId::from("TEST-001"),
2677 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
2678 position_id: PositionId::test_default(),
2679 account_id: AccountId::from("ACC-001"),
2680 opening_order_id: ClientOrderId::from("O-001"),
2681 entry: OrderSide::Buy,
2682 side: PositionSide::Long,
2683 signed_qty: 1.0,
2684 quantity: Quantity::default(),
2685 last_qty: Quantity::default(),
2686 last_px: Price::default(),
2687 currency: Currency::from("USD"),
2688 avg_px_open: 0.0,
2689 event_id: UUID4::default(),
2690 ts_event: UnixNanos::default(),
2691 ts_init: UnixNanos::default(),
2692 })
2693 }
2694
2695 fn make_position_changed() -> PositionEvent {
2696 let currency = Currency::from("USD");
2697 PositionEvent::PositionChanged(PositionChanged {
2698 trader_id: TraderId::from("TRADER-001"),
2699 strategy_id: StrategyId::from("TEST-001"),
2700 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
2701 position_id: PositionId::test_default(),
2702 account_id: AccountId::from("ACC-001"),
2703 opening_order_id: ClientOrderId::from("O-001"),
2704 entry: OrderSide::Buy,
2705 side: PositionSide::Long,
2706 signed_qty: 2.0,
2707 quantity: Quantity::default(),
2708 peak_quantity: Quantity::default(),
2709 last_qty: Quantity::default(),
2710 last_px: Price::default(),
2711 currency,
2712 avg_px_open: 0.0,
2713 avg_px_close: None,
2714 realized_return: 0.0,
2715 realized_pnl: None,
2716 unrealized_pnl: Money::zero(currency),
2717 event_id: UUID4::default(),
2718 ts_opened: UnixNanos::default(),
2719 ts_event: UnixNanos::default(),
2720 ts_init: UnixNanos::default(),
2721 })
2722 }
2723
2724 fn make_position_closed() -> PositionEvent {
2725 let currency = Currency::from("USD");
2726 PositionEvent::PositionClosed(PositionClosed {
2727 trader_id: TraderId::from("TRADER-001"),
2728 strategy_id: StrategyId::from("TEST-001"),
2729 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
2730 position_id: PositionId::test_default(),
2731 account_id: AccountId::from("ACC-001"),
2732 opening_order_id: ClientOrderId::from("O-001"),
2733 closing_order_id: Some(ClientOrderId::from("O-002")),
2734 entry: OrderSide::Buy,
2735 side: PositionSide::Flat,
2736 signed_qty: 0.0,
2737 quantity: Quantity::default(),
2738 peak_quantity: Quantity::default(),
2739 last_qty: Quantity::default(),
2740 last_px: Price::default(),
2741 currency,
2742 avg_px_open: 0.0,
2743 avg_px_close: None,
2744 realized_return: 0.0,
2745 realized_pnl: None,
2746 unrealized_pnl: Money::zero(currency),
2747 duration: 0,
2748 event_id: UUID4::default(),
2749 ts_opened: UnixNanos::default(),
2750 ts_closed: None,
2751 ts_event: UnixNanos::default(),
2752 ts_init: UnixNanos::default(),
2753 })
2754 }
2755
2756 fn make_position_adjusted() -> PositionEvent {
2757 PositionEvent::PositionAdjusted(PositionAdjusted {
2758 trader_id: TraderId::from("TRADER-001"),
2759 strategy_id: StrategyId::from("TEST-001"),
2760 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
2761 position_id: PositionId::test_default(),
2762 account_id: AccountId::from("ACC-001"),
2763 adjustment_type: PositionAdjustmentType::Funding,
2764 quantity_change: None,
2765 pnl_change: None,
2766 reason: None,
2767 event_id: UUID4::default(),
2768 ts_event: UnixNanos::default(),
2769 ts_init: UnixNanos::default(),
2770 })
2771 }
2772
2773 #[rstest]
2774 fn test_strategy_creation() {
2775 let strategy = create_test_strategy();
2776 assert_eq!(strategy.strategy_id(), Some(StrategyId::from("TEST-001")));
2777 assert!(!strategy.on_order_rejected_called);
2778 assert!(!strategy.on_position_opened_called);
2779 }
2780
2781 #[rstest]
2782 fn test_strategy_registration() {
2783 let mut strategy = create_test_strategy();
2784 register_strategy(&mut strategy);
2785
2786 assert!(strategy.is_registered());
2787 let _ = strategy.order().generate_client_order_id();
2788 let _ = strategy.portfolio().is_initialized();
2789 }
2790
2791 #[rstest]
2792 fn test_strategy_native_methods_are_available_on_strategy_type() {
2793 let mut strategy = create_test_strategy();
2794 register_strategy(&mut strategy);
2795
2796 drop(strategy.order_factory());
2797
2798 assert!(Rc::ptr_eq(
2799 &strategy.order_factory_rc(),
2800 strategy.core.order_factory.as_ref().unwrap()
2801 ));
2802 assert!(Rc::ptr_eq(
2803 &strategy.portfolio_rc(),
2804 strategy.core.portfolio.as_ref().unwrap()
2805 ));
2806 }
2807
2808 #[rstest]
2809 fn test_handle_order_event_dispatches_to_handler() {
2810 let mut strategy = create_test_strategy();
2811 register_strategy(&mut strategy);
2812 start_strategy(&mut strategy);
2813
2814 let event = make_rejected(ClientOrderId::from("O-001"));
2815
2816 strategy.handle_order_event(event);
2817
2818 assert!(strategy.on_order_rejected_called);
2819 assert!(strategy.on_order_event_called);
2820 }
2821
2822 #[rstest]
2823 fn test_handle_order_fill_voided_dispatches_to_specific_handler() {
2824 let mut strategy = create_test_strategy();
2825 register_strategy(&mut strategy);
2826 start_strategy(&mut strategy);
2827
2828 strategy.handle_order_event(make_fill_voided(ClientOrderId::from("O-001"), false));
2829
2830 assert!(strategy.on_order_fill_voided_called);
2831 assert!(strategy.on_order_event_called);
2832 }
2833
2834 #[rstest]
2835 #[case::opened(make_position_opened())]
2836 #[case::changed(make_position_changed())]
2837 #[case::closed(make_position_closed())]
2838 fn test_handle_position_event_dispatches_to_handler(#[case] event: PositionEvent) {
2839 let mut strategy = create_test_strategy();
2840 register_strategy(&mut strategy);
2841 start_strategy(&mut strategy);
2842
2843 let expected_opened = matches!(event, PositionEvent::PositionOpened(_));
2844 let expected_changed = matches!(event, PositionEvent::PositionChanged(_));
2845 let expected_closed = matches!(event, PositionEvent::PositionClosed(_));
2846
2847 strategy.handle_position_event(event);
2848
2849 assert_eq!(strategy.on_position_opened_called, expected_opened);
2850 assert_eq!(strategy.on_position_changed_called, expected_changed);
2851 assert_eq!(strategy.on_position_closed_called, expected_closed);
2852 assert!(strategy.on_position_event_called);
2853 }
2854
2855 #[rstest]
2856 fn test_handle_position_event_skips_dispatch_when_stopped() {
2857 let mut strategy = create_test_strategy();
2858 register_strategy(&mut strategy);
2859 start_strategy(&mut strategy);
2860 stop_strategy(&mut strategy);
2861 assert_eq!(strategy.state(), ComponentState::Stopped);
2862
2863 strategy.handle_position_event(make_position_opened());
2864
2865 assert!(!strategy.on_position_event_called);
2866 assert!(!strategy.on_position_opened_called);
2867 }
2868
2869 #[rstest]
2870 fn test_handle_position_event_skips_dispatch_for_adjusted() {
2871 let mut strategy = create_test_strategy();
2872 register_strategy(&mut strategy);
2873 start_strategy(&mut strategy);
2874
2875 strategy.handle_position_event(make_position_adjusted());
2876
2877 assert!(!strategy.on_position_event_called);
2878 assert!(!strategy.on_position_opened_called);
2879 assert!(!strategy.on_position_changed_called);
2880 assert!(!strategy.on_position_closed_called);
2881 }
2882
2883 #[rstest]
2884 fn test_strategy_default_handlers_do_not_panic() {
2885 let mut strategy = create_test_strategy();
2886
2887 strategy.on_order_initialized(OrderInitialized::default());
2888 strategy.on_order_event(OrderEventAny::Accepted(OrderAccepted::default()));
2889 strategy.on_order_denied(OrderDenied::default());
2890 strategy.on_order_emulated(OrderEmulated::default());
2891 strategy.on_order_released(OrderReleased::default());
2892 strategy.on_order_submitted(OrderSubmitted::default());
2893 strategy.on_order_rejected(OrderRejected::default());
2894 strategy.on_order_canceled(&OrderCanceled::default());
2895 strategy.on_order_expired(OrderExpired::default());
2896 strategy.on_order_triggered(OrderTriggered::default());
2897 strategy.on_order_pending_update(OrderPendingUpdate::default());
2898 strategy.on_order_pending_cancel(OrderPendingCancel::default());
2899 strategy.on_order_modify_rejected(OrderModifyRejected::default());
2900 strategy.on_order_cancel_rejected(OrderCancelRejected::default());
2901 strategy.on_order_updated(OrderUpdated::default());
2902 strategy.on_order_filled(&OrderFilledSpec::builder().build());
2903 strategy.on_order_fill_voided(&OrderFillVoidedSpec::builder().build());
2904 strategy.on_position_event(make_position_opened());
2905 }
2906
2907 #[rstest]
2908 fn test_submit_order_publishes_order_initialized_after_cache_insert_before_send() {
2909 let mut strategy = create_test_strategy();
2910 register_strategy(&mut strategy);
2911
2912 let order = make_initialized_market_order("O-20250208-INIT-001");
2913 let client_order_id = order.client_order_id();
2914 let cache_rc = strategy.core.cache_rc();
2915 let timeline = Rc::new(RefCell::new(Vec::new()));
2916 let event_messages = Rc::new(RefCell::new(Vec::new()));
2917
2918 let event_handler = {
2919 let event_messages = event_messages.clone();
2920 let timeline = timeline.clone();
2921 TypedHandler::from_with_id("events.order.initialized", move |event: &OrderEventAny| {
2922 assert!(cache_rc.borrow().order_exists(&client_order_id));
2923 assert!(matches!(event, OrderEventAny::Initialized(_)));
2924 event_messages.borrow_mut().push(event.clone());
2925 timeline.borrow_mut().push("init");
2926 })
2927 };
2928 let risk_handler = {
2929 let timeline = timeline.clone();
2930 TypedIntoHandler::from_with_id(
2931 "RiskEngine.queue_execute",
2932 move |command: TradingCommand| {
2933 assert!(matches!(command, TradingCommand::SubmitOrder(_)));
2934 timeline.borrow_mut().push("command");
2935 },
2936 )
2937 };
2938 msgbus::register_trading_command_endpoint(
2939 MessagingSwitchboard::risk_engine_queue_execute(),
2940 risk_handler,
2941 );
2942
2943 let topic = format!("events.order.{}", order.strategy_id());
2944 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
2945
2946 strategy
2947 .submit_order(order.clone(), None, None, None)
2948 .unwrap();
2949
2950 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
2951
2952 let event_messages = event_messages.borrow();
2953 assert_eq!(event_messages.len(), 1);
2954 assert_eq!(
2955 event_messages[0],
2956 OrderEventAny::Initialized(order.init_event().clone())
2957 );
2958 assert_eq!(timeline.borrow().as_slice(), &["init", "command"]);
2959 }
2960
2961 #[rstest]
2962 fn test_submit_order_routes_emulated_order_to_order_emulator() {
2963 let mut strategy = create_test_strategy();
2964 register_strategy(&mut strategy);
2965 let (emulator_handler, emulator_messages): (
2966 _,
2967 TypedIntoMessageSavingHandler<TradingCommand>,
2968 ) = get_typed_into_message_saving_handler(Some(Ustr::from("OrderEmulator.execute")));
2969 msgbus::register_trading_command_endpoint(
2970 MessagingSwitchboard::order_emulator_execute(),
2971 emulator_handler,
2972 );
2973 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
2974 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
2975 msgbus::register_trading_command_endpoint(
2976 MessagingSwitchboard::risk_engine_queue_execute(),
2977 risk_handler,
2978 );
2979 let order = OrderTestBuilder::new(OrderType::StopMarket)
2980 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
2981 .client_order_id(ClientOrderId::from("O-20250208-EMULATED-001"))
2982 .side(OrderSide::Buy)
2983 .trigger_price(Price::from("51000.0"))
2984 .quantity(Quantity::from(100_000))
2985 .emulation_trigger(TriggerType::BidAsk)
2986 .build();
2987 let client_order_id = order.client_order_id();
2988
2989 strategy.submit_order(order, None, None, None).unwrap();
2990
2991 let emulator_messages = emulator_messages.get_messages();
2992 assert_eq!(emulator_messages.len(), 1);
2993 assert!(matches!(
2994 emulator_messages.first(),
2995 Some(TradingCommand::SubmitOrder(command))
2996 if command.client_order_id == client_order_id
2997 ));
2998 assert!(risk_messages.get_messages().is_empty());
2999 }
3000
3001 #[rstest]
3002 fn test_submit_order_errors_when_strategy_not_registered() {
3003 let mut strategy = create_test_strategy();
3004 let order = make_initialized_market_order("O-20250208-UNREGISTERED-001");
3005
3006 let err = strategy
3007 .submit_order(order, None, None, None)
3008 .unwrap_err()
3009 .to_string();
3010
3011 assert_eq!(err, "Strategy not registered: trader_id is not set");
3012 }
3013
3014 #[rstest]
3015 fn test_required_account_id_errors_when_missing_for_strategy_event() {
3016 let order = make_initialized_market_order("O-20250208-NO-ACCOUNT-001");
3017
3018 let err = required_account_id(&order, "pending cancel")
3019 .unwrap_err()
3020 .to_string();
3021
3022 assert_eq!(
3023 err,
3024 "Cannot generate pending cancel event for O-20250208-NO-ACCOUNT-001: \
3025 account_id is not set"
3026 );
3027 }
3028
3029 #[rstest]
3030 fn test_submit_order_rejects_non_initialized_without_events() {
3031 let mut strategy = create_test_strategy();
3032 register_strategy(&mut strategy);
3033
3034 let order = make_accepted_market_order("O-20250208-ACCEPTED-001");
3035 let topic = format!("events.order.{}", order.strategy_id());
3036 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
3037 get_typed_message_saving_handler(Some(Ustr::from("events.order.invalid")));
3038
3039 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3040 let result = strategy.submit_order(order, None, None, None);
3041
3042 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3043
3044 assert!(result.is_err());
3045 assert!(
3046 result
3047 .unwrap_err()
3048 .to_string()
3049 .contains("expected INITIALIZED")
3050 );
3051 assert!(event_messages.get_messages().is_empty());
3052 }
3053
3054 #[rstest]
3055 fn test_submit_order_returns_error_when_cache_already_borrowed() {
3056 let mut strategy = create_test_strategy();
3057 register_strategy(&mut strategy);
3058
3059 let order = make_initialized_market_order("O-20250208-BORROWED-001");
3060 let cache_rc = strategy.core.cache_rc();
3061 let _cache = cache_rc.borrow();
3062
3063 let result = catch_unwind(AssertUnwindSafe(|| {
3064 strategy.submit_order(order, None, None, None)
3065 }));
3066
3067 let err = result
3068 .expect("submit_order should not panic")
3069 .unwrap_err()
3070 .to_string();
3071
3072 assert_eq!(
3073 err,
3074 "Cannot submit order O-20250208-BORROWED-001: cache is currently borrowed"
3075 );
3076 }
3077
3078 #[rstest]
3079 fn test_submit_order_list_publishes_order_initialized_after_cache_insert_before_send() {
3080 let mut strategy = create_test_strategy();
3081 register_strategy(&mut strategy);
3082
3083 let order_list_id = OrderListId::from("OL-20250208-LIST-INIT");
3084 let mut orders = vec![
3085 make_initialized_market_order("O-20250208-LIST-INIT-001"),
3086 make_initialized_market_order("O-20250208-LIST-INIT-002"),
3087 ];
3088
3089 for order in &mut orders {
3090 order.set_order_list_id(order_list_id);
3091 }
3092
3093 let client_order_id1 = orders[0].client_order_id();
3094 let client_order_id2 = orders[1].client_order_id();
3095 let cache_rc = strategy.core.cache_rc();
3096 let timeline = Rc::new(RefCell::new(Vec::new()));
3097 let event_messages = Rc::new(RefCell::new(Vec::new()));
3098
3099 let event_handler = {
3100 let event_messages = event_messages.clone();
3101 let timeline = timeline.clone();
3102 TypedHandler::from_with_id(
3103 "events.order.list_initialized",
3104 move |event: &OrderEventAny| {
3105 match event {
3106 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id1 => {
3107 let cache = cache_rc.borrow();
3108 assert!(cache.order_exists(&client_order_id1));
3109 assert!(cache.order_exists(&client_order_id2));
3110 assert!(cache.order_list_exists(&order_list_id));
3111 let order_list = cache.order_list(&order_list_id).unwrap();
3112 assert_eq!(
3113 order_list.client_order_ids.as_slice(),
3114 &[client_order_id1, client_order_id2]
3115 );
3116 timeline.borrow_mut().push("init1");
3117 }
3118 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id2 => {
3119 assert!(cache_rc.borrow().order_exists(&client_order_id2));
3120 timeline.borrow_mut().push("init2");
3121 }
3122 _ => panic!("unexpected order event {event:?}"),
3123 }
3124 event_messages.borrow_mut().push(event.clone());
3125 },
3126 )
3127 };
3128 let risk_handler = {
3129 let timeline = timeline.clone();
3130 TypedIntoHandler::from_with_id(
3131 "RiskEngine.queue_execute",
3132 move |command: TradingCommand| {
3133 assert!(matches!(command, TradingCommand::SubmitOrderList(_)));
3134 timeline.borrow_mut().push("command");
3135 },
3136 )
3137 };
3138 msgbus::register_trading_command_endpoint(
3139 MessagingSwitchboard::risk_engine_queue_execute(),
3140 risk_handler,
3141 );
3142
3143 let topic = format!("events.order.{}", orders[0].strategy_id());
3144 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3145
3146 strategy
3147 .submit_order_list(orders.clone(), None, None, None)
3148 .unwrap();
3149
3150 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3151
3152 let event_messages = event_messages.borrow();
3153 assert_eq!(event_messages.len(), 2);
3154 assert_eq!(
3155 event_messages[0],
3156 OrderEventAny::Initialized(orders[0].init_event().clone())
3157 );
3158 assert_eq!(
3159 event_messages[1],
3160 OrderEventAny::Initialized(orders[1].init_event().clone())
3161 );
3162 assert_eq!(timeline.borrow().as_slice(), &["init1", "init2", "command"]);
3163 }
3164
3165 #[rstest]
3166 fn test_submit_order_list_returns_error_when_cache_already_borrowed() {
3167 let mut strategy = create_test_strategy();
3168 register_strategy(&mut strategy);
3169
3170 let order_list_id = OrderListId::from("OL-20250208-BORROWED");
3171 let mut orders = vec![
3172 make_initialized_market_order("O-20250208-LIST-BORROWED-001"),
3173 make_initialized_market_order("O-20250208-LIST-BORROWED-002"),
3174 ];
3175
3176 for order in &mut orders {
3177 order.set_order_list_id(order_list_id);
3178 }
3179
3180 let cache_rc = strategy.core.cache_rc();
3181 let _cache = cache_rc.borrow();
3182
3183 let result = catch_unwind(AssertUnwindSafe(|| {
3184 strategy.submit_order_list(orders, None, None, None)
3185 }));
3186
3187 let err = result
3188 .expect("submit_order_list should not panic")
3189 .unwrap_err()
3190 .to_string();
3191
3192 assert_eq!(
3193 err,
3194 "Cannot submit order list OL-20250208-BORROWED: cache is currently borrowed"
3195 );
3196 }
3197
3198 #[rstest]
3199 fn test_submit_order_list_create_list_branch_publishes_init_after_cache_insert() {
3200 let mut strategy = create_test_strategy();
3201 register_strategy(&mut strategy);
3202
3203 let orders = vec![
3204 make_initialized_market_order("O-20250208-LIST-CREATE-001"),
3205 make_initialized_market_order("O-20250208-LIST-CREATE-002"),
3206 ];
3207
3208 let client_order_id1 = orders[0].client_order_id();
3209 let client_order_id2 = orders[1].client_order_id();
3210 let cache_rc = strategy.core.cache_rc();
3211 let timeline = Rc::new(RefCell::new(Vec::new()));
3212 let event_messages = Rc::new(RefCell::new(Vec::new()));
3213
3214 let event_handler = {
3215 let event_messages = event_messages.clone();
3216 let timeline = timeline.clone();
3217 TypedHandler::from_with_id(
3218 "events.order.list_create_initialized",
3219 move |event: &OrderEventAny| {
3220 match event {
3221 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id1 => {
3222 let cache = cache_rc.borrow();
3223 let cached_order1 = cache.order(&client_order_id1).unwrap();
3224 let cached_order2 = cache.order(&client_order_id2).unwrap();
3225 let order_list_id = cached_order1.order_list_id().unwrap();
3226 assert_eq!(cached_order2.order_list_id(), Some(order_list_id));
3227 assert!(cache.order_list_exists(&order_list_id));
3228 let order_list = cache.order_list(&order_list_id).unwrap();
3229 assert_eq!(
3230 order_list.client_order_ids.as_slice(),
3231 &[client_order_id1, client_order_id2]
3232 );
3233 timeline.borrow_mut().push("init1");
3234 }
3235 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id2 => {
3236 assert!(cache_rc.borrow().order_exists(&client_order_id2));
3237 timeline.borrow_mut().push("init2");
3238 }
3239 _ => panic!("unexpected order event {event:?}"),
3240 }
3241 event_messages.borrow_mut().push(event.clone());
3242 },
3243 )
3244 };
3245 let risk_handler = {
3246 let timeline = timeline.clone();
3247 TypedIntoHandler::from_with_id(
3248 "RiskEngine.queue_execute",
3249 move |command: TradingCommand| {
3250 assert!(matches!(command, TradingCommand::SubmitOrderList(_)));
3251 timeline.borrow_mut().push("command");
3252 },
3253 )
3254 };
3255 msgbus::register_trading_command_endpoint(
3256 MessagingSwitchboard::risk_engine_queue_execute(),
3257 risk_handler,
3258 );
3259
3260 let topic = format!("events.order.{}", orders[0].strategy_id());
3261 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3262
3263 strategy
3264 .submit_order_list(orders.clone(), None, None, None)
3265 .unwrap();
3266
3267 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3268
3269 let event_messages = event_messages.borrow();
3270 assert_eq!(event_messages.len(), 2);
3271 assert_eq!(
3272 event_messages[0],
3273 OrderEventAny::Initialized(orders[0].init_event().clone())
3274 );
3275 assert_eq!(
3276 event_messages[1],
3277 OrderEventAny::Initialized(orders[1].init_event().clone())
3278 );
3279 assert_eq!(timeline.borrow().as_slice(), &["init1", "init2", "command"]);
3280
3281 let cache = strategy.cache();
3282 let cached_order1 = cache.order(&client_order_id1).unwrap();
3283 let cached_order2 = cache.order(&client_order_id2).unwrap();
3284 let order_list_id = cached_order1.order_list_id().unwrap();
3285 assert_eq!(cached_order2.order_list_id(), Some(order_list_id));
3286
3287 let order_list = cache.order_list(&order_list_id).unwrap();
3288 assert_eq!(
3289 order_list.client_order_ids.as_slice(),
3290 &[client_order_id1, client_order_id2]
3291 );
3292 }
3293
3294 #[rstest]
3295 fn test_submit_order_list_routes_optional_params_to_risk() {
3296 let mut strategy = create_test_strategy();
3297 register_strategy(&mut strategy);
3298
3299 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3300 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3301 msgbus::register_trading_command_endpoint(
3302 MessagingSwitchboard::risk_engine_queue_execute(),
3303 risk_handler,
3304 );
3305
3306 let no_params_orders = vec![
3307 make_initialized_market_order("O-20250208-LIST-001"),
3308 make_initialized_market_order("O-20250208-LIST-002"),
3309 ];
3310 strategy
3311 .submit_order_list(no_params_orders, None, None, None)
3312 .unwrap();
3313
3314 let mut params = Params::new();
3315 params.insert(
3316 "routing_hint".to_string(),
3317 Value::String("prefer_batch".to_string()),
3318 );
3319 let param_orders = vec![
3320 make_initialized_market_order("O-20250208-LIST-003"),
3321 make_initialized_market_order("O-20250208-LIST-004"),
3322 ];
3323 strategy
3324 .submit_order_list(param_orders, None, None, Some(params.clone()))
3325 .unwrap();
3326
3327 let risk_messages = risk_messages.get_messages();
3328 assert_eq!(risk_messages.len(), 2);
3329 let Some(TradingCommand::SubmitOrderList(no_params_command)) = risk_messages.first() else {
3330 panic!("expected SubmitOrderList command");
3331 };
3332 let Some(TradingCommand::SubmitOrderList(param_command)) = risk_messages.get(1) else {
3333 panic!("expected SubmitOrderList command");
3334 };
3335 assert!(no_params_command.params.is_none());
3336 assert_eq!(param_command.params.as_ref(), Some(¶ms));
3337 }
3338
3339 #[rstest]
3340 fn test_modify_order_routes_non_emulated_orders_to_risk() {
3341 let mut strategy = create_test_strategy();
3342 register_strategy(&mut strategy);
3343
3344 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3345 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3346 msgbus::register_trading_command_endpoint(
3347 MessagingSwitchboard::risk_engine_queue_execute(),
3348 risk_handler,
3349 );
3350
3351 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3352 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3353 msgbus::register_trading_command_endpoint(
3354 MessagingSwitchboard::exec_engine_queue_execute(),
3355 exec_handler,
3356 );
3357
3358 let order = OrderAny::Market(MarketOrder::new(
3359 TraderId::from("TRADER-001"),
3360 StrategyId::from("TEST-001"),
3361 InstrumentId::from("BTCUSDT.BINANCE"),
3362 ClientOrderId::from("O-20250208-0003"),
3363 OrderSide::Buy,
3364 Quantity::from(100_000),
3365 TimeInForce::Gtc,
3366 UUID4::new(),
3367 UnixNanos::default(),
3368 false,
3369 false,
3370 None,
3371 None,
3372 None,
3373 None,
3374 None,
3375 None,
3376 None,
3377 None,
3378 ));
3379 add_order_to_cache(&strategy, &order);
3380
3381 strategy
3382 .modify_order(
3383 order.client_order_id(),
3384 Some(Quantity::from(200_000)),
3385 None,
3386 None,
3387 None,
3388 None,
3389 )
3390 .unwrap();
3391
3392 let risk_messages = risk_messages.get_messages();
3393 let exec_messages = exec_messages.get_messages();
3394
3395 assert_eq!(risk_messages.len(), 1);
3396 assert!(matches!(
3397 risk_messages.first(),
3398 Some(TradingCommand::ModifyOrder(_))
3399 ));
3400 assert!(exec_messages.is_empty());
3401 }
3402
3403 #[rstest]
3404 fn test_modify_order_marks_order_pending_update_locally_before_send() {
3405 let mut strategy = create_test_strategy();
3406 register_strategy(&mut strategy);
3407
3408 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3409 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3410 msgbus::register_trading_command_endpoint(
3411 MessagingSwitchboard::risk_engine_queue_execute(),
3412 risk_handler,
3413 );
3414
3415 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
3416 get_typed_message_saving_handler(Some(Ustr::from("events.order.pending_update")));
3417 let order = make_accepted_limit_order("O-20250208-UPDATE-001");
3418 let topic = format!("events.order.{}", order.strategy_id());
3419 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3420 add_order_to_cache(&strategy, &order);
3421
3422 strategy
3423 .modify_order(
3424 order.client_order_id(),
3425 None,
3426 Some(Price::from("51000.0")),
3427 None,
3428 None,
3429 None,
3430 )
3431 .unwrap();
3432
3433 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3434
3435 let cache = strategy.cache();
3436 let cached_order = cache.order(&order.client_order_id()).unwrap();
3437 assert_eq!(cached_order.status(), OrderStatus::PendingUpdate);
3438
3439 let risk_messages = risk_messages.get_messages();
3440 assert_eq!(risk_messages.len(), 1);
3441 assert!(matches!(
3442 risk_messages.first(),
3443 Some(TradingCommand::ModifyOrder(_))
3444 ));
3445
3446 let event_messages = event_messages.get_messages();
3447 assert_eq!(event_messages.len(), 1);
3448 assert!(matches!(
3449 event_messages.first(),
3450 Some(OrderEventAny::PendingUpdate(_))
3451 ));
3452 }
3453
3454 #[rstest]
3455 fn test_modify_orders_marks_orders_pending_update_locally_before_send() {
3456 let mut strategy = create_test_strategy();
3457 register_strategy(&mut strategy);
3458
3459 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3460 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3461 msgbus::register_trading_command_endpoint(
3462 MessagingSwitchboard::risk_engine_queue_execute(),
3463 risk_handler,
3464 );
3465
3466 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
3467 get_typed_message_saving_handler(Some(Ustr::from("events.order.batch_pending_update")));
3468 let order1 = make_accepted_limit_order("O-20250208-BATCH-UPDATE-001");
3469 let order2 = make_accepted_limit_order("O-20250208-BATCH-UPDATE-002");
3470 let topic = format!("events.order.{}", order1.strategy_id());
3471 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3472 add_order_to_cache(&strategy, &order1);
3473 add_order_to_cache(&strategy, &order2);
3474
3475 strategy
3476 .modify_orders(
3477 vec![
3478 (
3479 order1.client_order_id(),
3480 None,
3481 Some(Price::from("51000.0")),
3482 None,
3483 ),
3484 (
3485 order2.client_order_id(),
3486 Some(Quantity::from("2.0")),
3487 None,
3488 None,
3489 ),
3490 ],
3491 None,
3492 None,
3493 )
3494 .unwrap();
3495
3496 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3497
3498 let cache = strategy.cache();
3499 let cached_order1 = cache.order(&order1.client_order_id()).unwrap();
3500 let cached_order2 = cache.order(&order2.client_order_id()).unwrap();
3501 assert_eq!(cached_order1.status(), OrderStatus::PendingUpdate);
3502 assert_eq!(cached_order2.status(), OrderStatus::PendingUpdate);
3503
3504 let risk_messages = risk_messages.get_messages();
3505 assert_eq!(risk_messages.len(), 1);
3506 let Some(TradingCommand::ModifyOrders(command)) = risk_messages.first() else {
3507 panic!("expected BatchModifyOrders command");
3508 };
3509 assert_eq!(command.modifies.len(), 2);
3510 assert_eq!(
3511 command
3512 .modifies
3513 .iter()
3514 .map(|modify| modify.client_order_id)
3515 .collect::<Vec<_>>(),
3516 vec![order1.client_order_id(), order2.client_order_id()]
3517 );
3518
3519 let event_messages = event_messages.get_messages();
3520 assert_eq!(event_messages.len(), 2);
3521 assert!(
3522 event_messages
3523 .iter()
3524 .all(|event| matches!(event, OrderEventAny::PendingUpdate(_)))
3525 );
3526 }
3527
3528 #[rstest]
3529 fn test_cancel_order_marks_order_pending_cancel_locally_before_send() {
3530 let mut strategy = create_test_strategy();
3531 register_strategy(&mut strategy);
3532
3533 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3534 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3535 msgbus::register_trading_command_endpoint(
3536 MessagingSwitchboard::exec_engine_queue_execute(),
3537 exec_handler,
3538 );
3539
3540 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
3541 get_typed_message_saving_handler(Some(Ustr::from("events.order.pending_cancel")));
3542 let order = make_accepted_market_order("O-20250208-CANCEL-001");
3543 let topic = format!("events.order.{}", order.strategy_id());
3544 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3545 add_order_to_cache(&strategy, &order);
3546
3547 strategy
3548 .cancel_order(order.client_order_id(), None, None)
3549 .unwrap();
3550
3551 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3552
3553 let cache = strategy.cache();
3554 let cached_order = cache.order(&order.client_order_id()).unwrap();
3555 assert_eq!(cached_order.status(), OrderStatus::PendingCancel);
3556 let cache = strategy.core.cache_ref();
3557 assert!(cache.is_order_pending_cancel_local(&order.client_order_id()));
3558
3559 let exec_messages = exec_messages.get_messages();
3560 assert_eq!(exec_messages.len(), 1);
3561 assert!(matches!(
3562 exec_messages.first(),
3563 Some(TradingCommand::CancelOrder(_))
3564 ));
3565
3566 let event_messages = event_messages.get_messages();
3567 assert_eq!(event_messages.len(), 1);
3568 assert!(matches!(
3569 event_messages.first(),
3570 Some(OrderEventAny::PendingCancel(_))
3571 ));
3572 }
3573
3574 #[rstest]
3575 fn test_cancel_orders_marks_orders_pending_cancel_locally_before_send() {
3576 let mut strategy = create_test_strategy();
3577 register_strategy(&mut strategy);
3578
3579 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3580 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3581 msgbus::register_trading_command_endpoint(
3582 MessagingSwitchboard::exec_engine_queue_execute(),
3583 exec_handler,
3584 );
3585
3586 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
3587 get_typed_message_saving_handler(Some(Ustr::from("events.order.batch_pending_cancel")));
3588 let order1 = make_accepted_market_order("O-20250208-CANCEL-001");
3589 let order2 = make_accepted_market_order("O-20250208-CANCEL-002");
3590 let topic = format!("events.order.{}", order1.strategy_id());
3591 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
3592 add_order_to_cache(&strategy, &order1);
3593 add_order_to_cache(&strategy, &order2);
3594
3595 strategy
3596 .cancel_orders(
3597 vec![order1.client_order_id(), order2.client_order_id()],
3598 None,
3599 None,
3600 )
3601 .unwrap();
3602
3603 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
3604
3605 let cache = strategy.cache();
3606 let cached_order1 = cache.order(&order1.client_order_id()).unwrap();
3607 let cached_order2 = cache.order(&order2.client_order_id()).unwrap();
3608 assert_eq!(cached_order1.status(), OrderStatus::PendingCancel);
3609 assert_eq!(cached_order2.status(), OrderStatus::PendingCancel);
3610 let cache = strategy.core.cache_ref();
3611 assert!(cache.is_order_pending_cancel_local(&order1.client_order_id()));
3612 assert!(cache.is_order_pending_cancel_local(&order2.client_order_id()));
3613
3614 let exec_messages = exec_messages.get_messages();
3615 assert_eq!(exec_messages.len(), 1);
3616 let Some(TradingCommand::CancelOrders(command)) = exec_messages.first() else {
3617 panic!("expected BatchCancelOrders command");
3618 };
3619 assert_eq!(command.cancels.len(), 2);
3620
3621 let event_messages = event_messages.get_messages();
3622 assert_eq!(event_messages.len(), 2);
3623 assert!(
3624 event_messages
3625 .iter()
3626 .all(|event| matches!(event, OrderEventAny::PendingCancel(_)))
3627 );
3628 }
3629
3630 #[rstest]
3631 fn test_cancel_order_updates_own_book_status_before_send() {
3632 let mut strategy = create_test_strategy();
3633 register_strategy(&mut strategy);
3634
3635 let (exec_handler, _exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3636 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3637 msgbus::register_trading_command_endpoint(
3638 MessagingSwitchboard::exec_engine_queue_execute(),
3639 exec_handler,
3640 );
3641
3642 let order = make_accepted_limit_order("O-20250208-CANCEL-OWN-BOOK-001");
3643 add_order_to_cache_and_own_book(&strategy, &order);
3644
3645 strategy
3646 .cancel_order(order.client_order_id(), None, None)
3647 .unwrap();
3648
3649 let mut accepted = AHashSet::new();
3650 accepted.insert(OrderStatus::Accepted);
3651 let mut pending_cancel = AHashSet::new();
3652 pending_cancel.insert(OrderStatus::PendingCancel);
3653
3654 let cache = strategy.cache();
3655 let own_book = cache.own_order_book(&order.instrument_id()).unwrap();
3656 assert!(own_book.bids_as_map(Some(&accepted), None, None).is_empty());
3657 let pending_bids = own_book.bids_as_map(Some(&pending_cancel), None, None);
3658 assert_eq!(pending_bids.values().map(Vec::len).sum::<usize>(), 1);
3659 }
3660
3661 #[rstest]
3662 fn test_cancel_order_returns_error_when_not_in_cache() {
3663 let mut strategy = create_test_strategy();
3664 register_strategy(&mut strategy);
3665
3666 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3667 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3668 msgbus::register_trading_command_endpoint(
3669 MessagingSwitchboard::exec_engine_queue_execute(),
3670 exec_handler,
3671 );
3672
3673 let missing_id = ClientOrderId::from("O-MISSING");
3674 let err = strategy
3675 .cancel_order(missing_id, None, None)
3676 .expect_err("expected cancel_order to fail when order is not in cache");
3677
3678 assert_eq!(
3679 err.to_string(),
3680 format!("Cannot cancel order: {ORDER_NOT_FOUND}: {missing_id}")
3681 );
3682 assert!(exec_messages.get_messages().is_empty());
3683 }
3684
3685 #[rstest]
3686 fn test_modify_order_returns_error_when_not_in_cache() {
3687 let mut strategy = create_test_strategy();
3688 register_strategy(&mut strategy);
3689
3690 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3691 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3692 msgbus::register_trading_command_endpoint(
3693 MessagingSwitchboard::risk_engine_queue_execute(),
3694 risk_handler,
3695 );
3696
3697 let missing_id = ClientOrderId::from("O-MISSING");
3698 let err = strategy
3699 .modify_order(missing_id, Some(Quantity::from(1)), None, None, None, None)
3700 .expect_err("expected modify_order to fail when order is not in cache");
3701
3702 assert_eq!(
3703 err.to_string(),
3704 format!("Cannot modify order: {ORDER_NOT_FOUND}: {missing_id}")
3705 );
3706 assert!(risk_messages.get_messages().is_empty());
3707 }
3708
3709 #[rstest]
3710 fn test_modify_orders_returns_error_when_any_id_missing() {
3711 let mut strategy = create_test_strategy();
3712 register_strategy(&mut strategy);
3713
3714 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3715 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
3716 msgbus::register_trading_command_endpoint(
3717 MessagingSwitchboard::risk_engine_queue_execute(),
3718 risk_handler,
3719 );
3720
3721 let order = make_accepted_limit_order("O-PRESENT");
3722 add_order_to_cache(&strategy, &order);
3723
3724 let missing_id = ClientOrderId::from("O-MISSING");
3725 let err = strategy
3726 .modify_orders(
3727 vec![
3728 (
3729 order.client_order_id(),
3730 None,
3731 Some(Price::from("51000.0")),
3732 None,
3733 ),
3734 (missing_id, Some(Quantity::from("2.0")), None, None),
3735 ],
3736 None,
3737 None,
3738 )
3739 .expect_err("expected modify_orders to fail when any id is missing");
3740
3741 assert_eq!(
3742 err.to_string(),
3743 format!("Cannot modify order: {ORDER_NOT_FOUND}: {missing_id}")
3744 );
3745 assert!(risk_messages.get_messages().is_empty());
3746 }
3747
3748 #[rstest]
3749 fn test_cancel_orders_returns_error_when_any_id_missing() {
3750 let mut strategy = create_test_strategy();
3751 register_strategy(&mut strategy);
3752
3753 let (exec_handler, exec_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
3754 get_typed_into_message_saving_handler(Some(Ustr::from("ExecEngine.queue_execute")));
3755 msgbus::register_trading_command_endpoint(
3756 MessagingSwitchboard::exec_engine_queue_execute(),
3757 exec_handler,
3758 );
3759
3760 let order = make_accepted_limit_order("O-PRESENT");
3761 add_order_to_cache(&strategy, &order);
3762
3763 let missing_id = ClientOrderId::from("O-MISSING");
3764 let err = strategy
3765 .cancel_orders(vec![order.client_order_id(), missing_id], None, None)
3766 .expect_err("expected cancel_orders to fail when any id is missing");
3767
3768 assert_eq!(
3769 err.to_string(),
3770 format!("Cannot cancel order: {ORDER_NOT_FOUND}: {missing_id}")
3771 );
3772 assert!(exec_messages.get_messages().is_empty());
3773 }
3774
3775 #[rstest]
3778 fn test_has_gtd_expiry_timer_when_timer_not_set() {
3779 let mut strategy = create_test_strategy();
3780 let client_order_id = ClientOrderId::from("O-001");
3781
3782 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
3783 }
3784
3785 #[rstest]
3786 fn test_has_gtd_expiry_timer_when_timer_set() {
3787 let mut strategy = create_test_strategy();
3788 let client_order_id = ClientOrderId::from("O-001");
3789
3790 strategy
3791 .core
3792 .gtd_timers
3793 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
3794
3795 assert!(strategy.has_gtd_expiry_timer(&client_order_id));
3796 }
3797
3798 #[rstest]
3799 fn test_cancel_gtd_expiry_removes_timer() {
3800 let mut strategy = create_test_strategy();
3801 register_strategy(&mut strategy);
3802
3803 let client_order_id = ClientOrderId::from("O-001");
3804 strategy
3805 .core
3806 .gtd_timers
3807 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
3808
3809 strategy.cancel_gtd_expiry(&client_order_id);
3810
3811 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
3812 }
3813
3814 #[rstest]
3815 fn test_cancel_gtd_expiry_when_timer_not_set() {
3816 let mut strategy = create_test_strategy();
3817 register_strategy(&mut strategy);
3818
3819 let client_order_id = ClientOrderId::from("O-001");
3820
3821 strategy.cancel_gtd_expiry(&client_order_id);
3822
3823 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
3824 }
3825
3826 #[rstest]
3827 #[case::filled(make_filled)]
3828 #[case::canceled(make_canceled)]
3829 #[case::rejected(make_rejected)]
3830 #[case::expired(make_expired)]
3831 #[case::fill_voided(make_terminal_fill_voided)]
3832 fn test_handle_order_event_cancels_gtd_timer_for_terminal_event(
3833 #[case] make_event: fn(ClientOrderId) -> OrderEventAny,
3834 ) {
3835 let mut strategy = create_test_strategy();
3836 register_strategy(&mut strategy);
3837 start_strategy(&mut strategy);
3838
3839 let client_order_id = ClientOrderId::from("O-001");
3840 strategy
3841 .core
3842 .gtd_timers
3843 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
3844
3845 strategy.handle_order_event(make_event(client_order_id));
3846
3847 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
3848 }
3849
3850 #[rstest]
3851 #[case::partial_fill(make_filled)]
3852 #[case::non_reopened_fill_void(make_terminal_fill_voided)]
3853 fn test_handle_order_event_keeps_gtd_timer_when_cached_order_remains_open(
3854 #[case] make_event: fn(ClientOrderId) -> OrderEventAny,
3855 ) {
3856 let mut strategy = create_test_strategy();
3857 register_strategy(&mut strategy);
3858 start_strategy(&mut strategy);
3859
3860 let client_order_id = ClientOrderId::from("O-001");
3861 let order = make_accepted_limit_order(client_order_id.as_str());
3862 add_order_to_cache(&strategy, &order);
3863 strategy
3864 .core
3865 .gtd_timers
3866 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
3867
3868 strategy.handle_order_event(make_event(client_order_id));
3869
3870 assert!(strategy.has_gtd_expiry_timer(&client_order_id));
3871 }
3872
3873 #[rstest]
3874 #[case::filled(make_filled)]
3875 #[case::canceled(make_canceled)]
3876 #[case::rejected(make_rejected)]
3877 #[case::expired(make_expired)]
3878 #[case::fill_voided(make_terminal_fill_voided)]
3879 fn test_handle_order_event_cancels_gtd_timer_when_stopped(
3880 #[case] make_event: fn(ClientOrderId) -> OrderEventAny,
3881 ) {
3882 let mut strategy = create_test_strategy();
3883 register_strategy(&mut strategy);
3884 start_strategy(&mut strategy);
3885
3886 let client_order_id = ClientOrderId::from("O-001");
3887 strategy
3888 .core
3889 .gtd_timers
3890 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
3891
3892 stop_strategy(&mut strategy);
3893 assert_eq!(strategy.state(), ComponentState::Stopped);
3894
3895 strategy.handle_order_event(make_event(client_order_id));
3896
3897 assert!(!strategy.has_gtd_expiry_timer(&client_order_id));
3898 }
3899
3900 #[rstest]
3901 fn test_handle_order_event_skips_gtd_cancel_for_non_terminal() {
3902 let mut strategy = create_test_strategy();
3903 register_strategy(&mut strategy);
3904 start_strategy(&mut strategy);
3905
3906 let client_order_id = ClientOrderId::from("O-001");
3907 strategy
3908 .core
3909 .gtd_timers
3910 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
3911
3912 strategy.handle_order_event(make_accepted(client_order_id));
3913
3914 assert!(strategy.has_gtd_expiry_timer(&client_order_id));
3915 }
3916
3917 #[rstest]
3918 fn test_handle_reopened_fill_void_keeps_gtd_timer() {
3919 let mut strategy = create_test_strategy();
3920 register_strategy(&mut strategy);
3921 start_strategy(&mut strategy);
3922
3923 let client_order_id = ClientOrderId::from("O-001");
3924 strategy
3925 .core
3926 .gtd_timers
3927 .insert(client_order_id, Ustr::from("GTD-EXPIRY:O-001"));
3928
3929 strategy.handle_order_event(make_fill_voided(client_order_id, true));
3930
3931 assert!(strategy.has_gtd_expiry_timer(&client_order_id));
3932 }
3933
3934 #[rstest]
3935 fn test_handle_order_event_skips_dispatch_when_stopped() {
3936 let mut strategy = create_test_strategy();
3937 register_strategy(&mut strategy);
3938 start_strategy(&mut strategy);
3939 stop_strategy(&mut strategy);
3940 assert_eq!(strategy.state(), ComponentState::Stopped);
3941
3942 strategy.handle_order_event(make_rejected(ClientOrderId::from("O-001")));
3943
3944 assert!(!strategy.on_order_event_called);
3945 assert!(!strategy.on_order_rejected_called);
3946 }
3947
3948 #[rstest]
3949 fn test_on_start_calls_reactivate_gtd_timers_when_enabled() {
3950 let config = StrategyConfig {
3951 strategy_id: Some(StrategyId::from("TEST-001")),
3952 order_id_tag: Some("001".to_string()),
3953 manage_gtd_expiry: true,
3954 ..Default::default()
3955 };
3956 let mut strategy = TestStrategy::new(config);
3957 register_strategy(&mut strategy);
3958
3959 let result = Strategy::on_start(&mut strategy);
3960 assert!(result.is_ok());
3961 }
3962
3963 #[rstest]
3964 fn test_on_start_does_not_panic_when_gtd_disabled() {
3965 let config = StrategyConfig {
3966 strategy_id: Some(StrategyId::from("TEST-001")),
3967 order_id_tag: Some("001".to_string()),
3968 manage_gtd_expiry: false,
3969 ..Default::default()
3970 };
3971 let mut strategy = TestStrategy::new(config);
3972 register_strategy(&mut strategy);
3973
3974 let result = Strategy::on_start(&mut strategy);
3975 assert!(result.is_ok());
3976 }
3977
3978 #[rstest]
3981 fn test_query_account_when_registered() {
3982 let mut strategy = create_test_strategy();
3983 register_strategy(&mut strategy);
3984
3985 let account_id = AccountId::from("ACC-001");
3986
3987 let result = strategy.query_account(account_id, None, None);
3988
3989 assert!(result.is_ok());
3990 }
3991
3992 #[rstest]
3993 fn test_query_account_with_client_id() {
3994 let mut strategy = create_test_strategy();
3995 register_strategy(&mut strategy);
3996
3997 let account_id = AccountId::from("ACC-001");
3998 let client_id = ClientId::from("BINANCE");
3999
4000 let result = strategy.query_account(account_id, Some(client_id), None);
4001
4002 assert!(result.is_ok());
4003 }
4004
4005 #[rstest]
4006 fn test_query_order_when_registered() {
4007 let mut strategy = create_test_strategy();
4008 register_strategy(&mut strategy);
4009
4010 let order = OrderAny::Market(MarketOrder::test_default());
4011
4012 let result = strategy.query_order(&order, None, None);
4013
4014 assert!(result.is_ok());
4015 }
4016
4017 #[rstest]
4018 fn test_query_order_with_client_id() {
4019 let mut strategy = create_test_strategy();
4020 register_strategy(&mut strategy);
4021
4022 let order = OrderAny::Market(MarketOrder::test_default());
4023 let client_id = ClientId::from("BINANCE");
4024
4025 let result = strategy.query_order(&order, Some(client_id), None);
4026
4027 assert!(result.is_ok());
4028 }
4029
4030 #[rstest]
4031 fn test_is_exiting_returns_false_by_default() {
4032 let strategy = create_test_strategy();
4033 assert!(!strategy.is_exiting());
4034 }
4035
4036 #[rstest]
4037 fn test_is_exiting_returns_true_when_set_manually() {
4038 let mut strategy = create_test_strategy();
4039 register_strategy(&mut strategy);
4040
4041 strategy.core.is_exiting = true;
4043
4044 assert!(strategy.is_exiting());
4045 }
4046
4047 #[rstest]
4048 fn test_market_exit_sets_is_exiting_flag() {
4049 let mut strategy = create_test_strategy();
4051 register_strategy(&mut strategy);
4052
4053 assert!(!strategy.core.is_exiting);
4054
4055 strategy.core.is_exiting = true;
4057 strategy.core.market_exit_attempts = 0;
4058
4059 assert!(strategy.core.is_exiting);
4060 assert_eq!(strategy.core.market_exit_attempts, 0);
4061 }
4062
4063 #[rstest]
4064 fn test_market_exit_uses_config_time_in_force_and_reduce_only() {
4065 let config = StrategyConfig {
4066 strategy_id: Some(StrategyId::from("TEST-001")),
4067 order_id_tag: Some("001".to_string()),
4068 market_exit_time_in_force: TimeInForce::Ioc,
4069 market_exit_reduce_only: false,
4070 ..Default::default()
4071 };
4072 let strategy = TestStrategy::new(config);
4073
4074 assert_eq!(
4075 strategy.core.config.market_exit_time_in_force,
4076 TimeInForce::Ioc
4077 );
4078 assert!(!strategy.core.config.market_exit_reduce_only);
4079 }
4080
4081 #[rstest]
4082 fn test_market_exit_resets_attempt_counter() {
4083 let mut strategy = create_test_strategy();
4084 register_strategy(&mut strategy);
4085
4086 strategy.core.market_exit_attempts = 50;
4088
4089 strategy.core.reset_market_exit_state();
4091
4092 assert_eq!(strategy.core.market_exit_attempts, 0);
4093 }
4094
4095 #[rstest]
4096 fn test_market_exit_second_call_returns_early_when_exiting() {
4097 let mut strategy = create_test_strategy();
4098 register_strategy(&mut strategy);
4099
4100 strategy.core.is_exiting = true;
4102
4103 let result = strategy.market_exit();
4105 assert!(result.is_ok());
4106 assert!(strategy.core.is_exiting);
4107 }
4108
4109 #[rstest]
4110 fn test_finalize_market_exit_resets_state() {
4111 let mut strategy = create_test_strategy();
4112 register_strategy(&mut strategy);
4113
4114 strategy.core.is_exiting = true;
4116 strategy.core.pending_stop = true;
4117 strategy.core.market_exit_attempts = 50;
4118
4119 strategy.finalize_market_exit();
4120
4121 assert!(!strategy.core.is_exiting);
4122 assert!(!strategy.core.pending_stop);
4123 assert_eq!(strategy.core.market_exit_attempts, 0);
4124 }
4125
4126 #[rstest]
4127 fn test_market_exit_config_defaults() {
4128 let config = StrategyConfig::default();
4129
4130 assert!(!config.manage_stop);
4131 assert_eq!(config.market_exit_interval_ms, 100);
4132 assert_eq!(config.market_exit_max_attempts, 100);
4133 }
4134
4135 #[rstest]
4136 fn test_market_exit_with_custom_config() {
4137 let config = StrategyConfig {
4138 strategy_id: Some(StrategyId::from("TEST-001")),
4139 manage_stop: true,
4140 market_exit_interval_ms: 50,
4141 market_exit_max_attempts: 200,
4142 ..Default::default()
4143 };
4144 let strategy = TestStrategy::new(config);
4145
4146 assert!(strategy.core.config.manage_stop);
4147 assert_eq!(strategy.core.config.market_exit_interval_ms, 50);
4148 assert_eq!(strategy.core.config.market_exit_max_attempts, 200);
4149 }
4150
4151 #[derive(Debug)]
4152 struct MarketExitHookTrackingStrategy {
4153 core: StrategyCore,
4154 on_market_exit_called: bool,
4155 post_market_exit_called: bool,
4156 }
4157
4158 impl MarketExitHookTrackingStrategy {
4159 fn new(config: StrategyConfig) -> Self {
4160 Self {
4161 core: StrategyCore::new(config),
4162 on_market_exit_called: false,
4163 post_market_exit_called: false,
4164 }
4165 }
4166 }
4167
4168 impl DataActor for MarketExitHookTrackingStrategy {}
4169
4170 nautilus_strategy!(MarketExitHookTrackingStrategy, {
4171 fn on_market_exit(&mut self) {
4172 self.on_market_exit_called = true;
4173 }
4174
4175 fn post_market_exit(&mut self) {
4176 self.post_market_exit_called = true;
4177 }
4178 });
4179
4180 #[rstest]
4181 fn test_market_exit_calls_on_market_exit_hook() {
4182 let config = StrategyConfig {
4183 strategy_id: Some(StrategyId::from("TEST-001")),
4184 order_id_tag: Some("001".to_string()),
4185 ..Default::default()
4186 };
4187 let mut strategy = MarketExitHookTrackingStrategy::new(config);
4188
4189 let trader_id = TraderId::from("TRADER-001");
4190 let clock = Rc::new(RefCell::new(TestClock::new()));
4191 let cache = Rc::new(RefCell::new(Cache::default()));
4192 let portfolio = Rc::new(RefCell::new(Portfolio::new(
4193 clock.clone(),
4194 cache.clone(),
4195 None,
4196 )));
4197 strategy
4198 .core
4199 .register(trader_id, clock, cache, portfolio)
4200 .unwrap();
4201 strategy.initialize().unwrap();
4202 strategy.start().unwrap();
4203
4204 let _ = strategy.market_exit();
4205
4206 assert!(strategy.on_market_exit_called);
4207 }
4208
4209 #[rstest]
4210 fn test_finalize_market_exit_calls_post_market_exit_hook() {
4211 let config = StrategyConfig {
4212 strategy_id: Some(StrategyId::from("TEST-001")),
4213 order_id_tag: Some("001".to_string()),
4214 ..Default::default()
4215 };
4216 let mut strategy = MarketExitHookTrackingStrategy::new(config);
4217
4218 let trader_id = TraderId::from("TRADER-001");
4219 let clock = Rc::new(RefCell::new(TestClock::new()));
4220 let cache = Rc::new(RefCell::new(Cache::default()));
4221 let portfolio = Rc::new(RefCell::new(Portfolio::new(
4222 clock.clone(),
4223 cache.clone(),
4224 None,
4225 )));
4226 strategy
4227 .core
4228 .register(trader_id, clock, cache, portfolio)
4229 .unwrap();
4230
4231 strategy.core.is_exiting = true;
4232 strategy.finalize_market_exit();
4233
4234 assert!(strategy.post_market_exit_called);
4235 }
4236
4237 #[derive(Debug)]
4238 struct FailingPostExitStrategy {
4239 core: StrategyCore,
4240 }
4241
4242 impl FailingPostExitStrategy {
4243 fn new(config: StrategyConfig) -> Self {
4244 Self {
4245 core: StrategyCore::new(config),
4246 }
4247 }
4248 }
4249
4250 impl DataActor for FailingPostExitStrategy {}
4251
4252 nautilus_strategy!(FailingPostExitStrategy, {
4253 fn post_market_exit(&mut self) {
4254 panic!("Simulated error in post_market_exit");
4255 }
4256 });
4257
4258 #[rstest]
4259 fn test_finalize_market_exit_handles_hook_panic() {
4260 let config = StrategyConfig {
4261 strategy_id: Some(StrategyId::from("TEST-001")),
4262 order_id_tag: Some("001".to_string()),
4263 ..Default::default()
4264 };
4265 let mut strategy = FailingPostExitStrategy::new(config);
4266
4267 let trader_id = TraderId::from("TRADER-001");
4268 let clock = Rc::new(RefCell::new(TestClock::new()));
4269 let cache = Rc::new(RefCell::new(Cache::default()));
4270 let portfolio = Rc::new(RefCell::new(Portfolio::new(
4271 clock.clone(),
4272 cache.clone(),
4273 None,
4274 )));
4275 strategy
4276 .core
4277 .register(trader_id, clock, cache, portfolio)
4278 .unwrap();
4279
4280 strategy.core.is_exiting = true;
4281 strategy.core.pending_stop = true;
4282
4283 strategy.finalize_market_exit();
4285
4286 assert!(!strategy.core.is_exiting);
4288 assert!(!strategy.core.pending_stop);
4289 }
4290
4291 #[rstest]
4292 fn test_check_market_exit_increments_attempts_before_finalizing() {
4293 let mut strategy = create_test_strategy();
4294 register_strategy(&mut strategy);
4295
4296 strategy.core.is_exiting = true;
4297 assert_eq!(strategy.core.market_exit_attempts, 0);
4298
4299 let event = TimeEvent::new(
4300 Ustr::from("MARKET_EXIT_CHECK:TEST-001"),
4301 UUID4::new(),
4302 UnixNanos::default(),
4303 UnixNanos::default(),
4304 );
4305 strategy.check_market_exit(event);
4306
4307 assert!(!strategy.core.is_exiting);
4311 assert_eq!(strategy.core.market_exit_attempts, 0);
4312 }
4313
4314 #[rstest]
4315 fn test_check_market_exit_finalizes_when_max_attempts_reached() {
4316 let config = StrategyConfig {
4317 strategy_id: Some(StrategyId::from("TEST-001")),
4318 order_id_tag: Some("001".to_string()),
4319 market_exit_max_attempts: 3,
4320 ..Default::default()
4321 };
4322 let mut strategy = TestStrategy::new(config);
4323 register_strategy(&mut strategy);
4324
4325 strategy.core.is_exiting = true;
4326 strategy.core.market_exit_attempts = 2; let event = TimeEvent::new(
4329 Ustr::from("MARKET_EXIT_CHECK:TEST-001"),
4330 UUID4::new(),
4331 UnixNanos::default(),
4332 UnixNanos::default(),
4333 );
4334 strategy.check_market_exit(event);
4335
4336 assert!(!strategy.core.is_exiting);
4338 assert_eq!(strategy.core.market_exit_attempts, 0);
4339 }
4340
4341 #[rstest]
4342 fn test_check_market_exit_finalizes_when_no_orders_or_positions() {
4343 let mut strategy = create_test_strategy();
4344 register_strategy(&mut strategy);
4345
4346 strategy.core.is_exiting = true;
4347
4348 let event = TimeEvent::new(
4349 Ustr::from("MARKET_EXIT_CHECK:TEST-001"),
4350 UUID4::new(),
4351 UnixNanos::default(),
4352 UnixNanos::default(),
4353 );
4354 strategy.check_market_exit(event);
4355
4356 assert!(!strategy.core.is_exiting);
4358 }
4359
4360 #[rstest]
4361 fn test_market_exit_timer_name_format() {
4362 let config = StrategyConfig {
4363 strategy_id: Some(StrategyId::from("MY-STRATEGY-001")),
4364 ..Default::default()
4365 };
4366 let strategy = TestStrategy::new(config);
4367
4368 assert_eq!(
4369 strategy.core.market_exit_timer_name.as_str(),
4370 "MARKET_EXIT_CHECK:MY-STRATEGY-001"
4371 );
4372 }
4373
4374 #[rstest]
4375 fn test_reset_market_exit_state() {
4376 let mut strategy = create_test_strategy();
4377
4378 strategy.core.is_exiting = true;
4379 strategy.core.pending_stop = true;
4380 strategy.core.market_exit_attempts = 50;
4381
4382 strategy.core.reset_market_exit_state();
4383
4384 assert!(!strategy.core.is_exiting);
4385 assert!(!strategy.core.pending_stop);
4386 assert_eq!(strategy.core.market_exit_attempts, 0);
4387 }
4388
4389 #[rstest]
4390 fn test_cancel_market_exit_resets_state_without_hooks() {
4391 let config = StrategyConfig {
4392 strategy_id: Some(StrategyId::from("TEST-001")),
4393 order_id_tag: Some("001".to_string()),
4394 ..Default::default()
4395 };
4396 let mut strategy = MarketExitHookTrackingStrategy::new(config);
4397
4398 let trader_id = TraderId::from("TRADER-001");
4399 let clock = Rc::new(RefCell::new(TestClock::new()));
4400 let cache = Rc::new(RefCell::new(Cache::default()));
4401 let portfolio = Rc::new(RefCell::new(Portfolio::new(
4402 clock.clone(),
4403 cache.clone(),
4404 None,
4405 )));
4406 strategy
4407 .core
4408 .register(trader_id, clock, cache, portfolio)
4409 .unwrap();
4410
4411 strategy.core.is_exiting = true;
4413 strategy.core.pending_stop = true;
4414 strategy.core.market_exit_attempts = 50;
4415
4416 strategy.cancel_market_exit();
4418
4419 assert!(!strategy.core.is_exiting);
4421 assert!(!strategy.core.pending_stop);
4422 assert_eq!(strategy.core.market_exit_attempts, 0);
4423
4424 assert!(!strategy.on_market_exit_called);
4426 assert!(!strategy.post_market_exit_called);
4427 }
4428
4429 #[rstest]
4430 fn test_market_exit_returns_early_when_not_running() {
4431 let mut strategy = create_test_strategy();
4432 register_strategy(&mut strategy);
4433
4434 assert!(!strategy.is_running());
4436
4437 let result = strategy.market_exit();
4438
4439 assert!(result.is_ok());
4441 assert!(!strategy.core.is_exiting);
4442 }
4443
4444 #[rstest]
4445 fn test_stop_with_manage_stop_false_cleans_up_active_exit() {
4446 let config = StrategyConfig {
4447 strategy_id: Some(StrategyId::from("TEST-001")),
4448 order_id_tag: Some("001".to_string()),
4449 manage_stop: false,
4450 ..Default::default()
4451 };
4452 let mut strategy = TestStrategy::new(config);
4453 register_strategy(&mut strategy);
4454
4455 strategy.core.is_exiting = true;
4457 strategy.core.market_exit_attempts = 5;
4458
4459 let should_proceed = Strategy::stop(&mut strategy);
4461
4462 assert!(should_proceed);
4464 assert!(!strategy.core.is_exiting);
4465 assert_eq!(strategy.core.market_exit_attempts, 0);
4466 }
4467
4468 #[rstest]
4469 fn test_stop_with_manage_stop_true_defers_when_running() {
4470 let config = StrategyConfig {
4471 strategy_id: Some(StrategyId::from("TEST-001")),
4472 order_id_tag: Some("001".to_string()),
4473 manage_stop: true,
4474 ..Default::default()
4475 };
4476 let mut strategy = TestStrategy::new(config);
4477
4478 let trader_id = TraderId::from("TRADER-001");
4480 let clock = Rc::new(RefCell::new(TestClock::new()));
4481 clock
4482 .borrow_mut()
4483 .register_default_handler(TimeEventCallback::from(|_event: TimeEvent| {}));
4484 let cache = Rc::new(RefCell::new(Cache::default()));
4485 let portfolio = Rc::new(RefCell::new(Portfolio::new(
4486 clock.clone(),
4487 cache.clone(),
4488 None,
4489 )));
4490 strategy
4491 .core
4492 .register(trader_id, clock, cache, portfolio)
4493 .unwrap();
4494 strategy.initialize().unwrap();
4495 strategy.start().unwrap();
4496
4497 let should_proceed = Strategy::stop(&mut strategy);
4498
4499 assert!(!should_proceed);
4501 assert!(strategy.core.pending_stop);
4502 }
4503
4504 #[rstest]
4505 fn test_stop_with_manage_stop_true_returns_early_if_pending() {
4506 let config = StrategyConfig {
4507 strategy_id: Some(StrategyId::from("TEST-001")),
4508 order_id_tag: Some("001".to_string()),
4509 manage_stop: true,
4510 ..Default::default()
4511 };
4512 let mut strategy = TestStrategy::new(config);
4513 register_strategy(&mut strategy);
4514 start_strategy(&mut strategy);
4515 strategy.core.pending_stop = true;
4516
4517 let should_proceed = Strategy::stop(&mut strategy);
4519
4520 assert!(!should_proceed);
4522 assert!(strategy.core.pending_stop);
4523 }
4524
4525 #[rstest]
4526 fn test_stop_with_manage_stop_true_proceeds_when_not_running() {
4527 let config = StrategyConfig {
4528 strategy_id: Some(StrategyId::from("TEST-001")),
4529 order_id_tag: Some("001".to_string()),
4530 manage_stop: true,
4531 ..Default::default()
4532 };
4533 let mut strategy = TestStrategy::new(config);
4534 register_strategy(&mut strategy);
4535
4536 assert!(!strategy.is_running());
4538
4539 let should_proceed = Strategy::stop(&mut strategy);
4540
4541 assert!(should_proceed);
4543 }
4544
4545 #[rstest]
4546 fn test_finalize_market_exit_stops_strategy_when_pending() {
4547 let config = StrategyConfig {
4548 strategy_id: Some(StrategyId::from("TEST-001")),
4549 order_id_tag: Some("001".to_string()),
4550 ..Default::default()
4551 };
4552 let mut strategy = TestStrategy::new(config);
4553 register_strategy(&mut strategy);
4554 start_strategy(&mut strategy);
4555
4556 strategy.core.is_exiting = true;
4558 strategy.core.pending_stop = true;
4559
4560 strategy.finalize_market_exit();
4561
4562 assert_eq!(strategy.state(), ComponentState::Stopped);
4564 assert!(!strategy.core.is_exiting);
4565 assert!(!strategy.core.pending_stop);
4566 }
4567
4568 #[rstest]
4569 fn test_finalize_market_exit_stays_running_when_not_pending() {
4570 let config = StrategyConfig {
4571 strategy_id: Some(StrategyId::from("TEST-001")),
4572 order_id_tag: Some("001".to_string()),
4573 ..Default::default()
4574 };
4575 let mut strategy = TestStrategy::new(config);
4576 register_strategy(&mut strategy);
4577 start_strategy(&mut strategy);
4578
4579 strategy.core.is_exiting = true;
4581 strategy.core.pending_stop = false;
4582
4583 strategy.finalize_market_exit();
4584
4585 assert_eq!(strategy.state(), ComponentState::Running);
4587 assert!(!strategy.core.is_exiting);
4588 }
4589
4590 #[rstest]
4591 fn test_submit_order_denied_during_market_exit_when_not_reduce_only() {
4592 let mut strategy = create_test_strategy();
4593 register_strategy(&mut strategy);
4594 start_strategy(&mut strategy);
4595 strategy.core.is_exiting = true;
4596
4597 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
4598 get_typed_message_saving_handler(Some(Ustr::from("events.order.denied")));
4599 let order = OrderAny::Market(MarketOrder::new(
4600 TraderId::from("TRADER-001"),
4601 StrategyId::from("TEST-001"),
4602 InstrumentId::from("BTCUSDT.BINANCE"),
4603 ClientOrderId::from("O-20250208-0001"),
4604 OrderSide::Buy,
4605 Quantity::from(100_000),
4606 TimeInForce::Gtc,
4607 UUID4::new(),
4608 UnixNanos::default(),
4609 false, false,
4611 None,
4612 None,
4613 None,
4614 None,
4615 None,
4616 None,
4617 None,
4618 None,
4619 ));
4620 let topic = format!("events.order.{}", order.strategy_id());
4621 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4622 let client_order_id = order.client_order_id();
4623 let result = strategy.submit_order(order.clone(), None, None, None);
4624
4625 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4626
4627 assert!(result.is_ok());
4628 let cache = strategy.cache();
4629 let cached_order = cache.order(&client_order_id).unwrap();
4630 assert_eq!(cached_order.status(), OrderStatus::Denied);
4631
4632 let event_messages = event_messages.get_messages();
4633 assert_eq!(event_messages.len(), 2);
4634 assert_eq!(
4635 event_messages[0],
4636 OrderEventAny::Initialized(order.init_event().clone())
4637 );
4638 let OrderEventAny::Denied(denied) = &event_messages[1] else {
4639 panic!("expected OrderDenied event");
4640 };
4641 assert_eq!(denied.reason, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
4642 }
4643
4644 #[rstest]
4645 fn test_submit_order_list_denied_during_market_exit_publishes_init_then_denied_events() {
4646 let mut strategy = create_test_strategy();
4647 register_strategy(&mut strategy);
4648 start_strategy(&mut strategy);
4649 strategy.core.is_exiting = true;
4650
4651 let orders = vec![
4652 make_initialized_market_order("O-20250208-LIST-DENY-001"),
4653 make_initialized_market_order("O-20250208-LIST-DENY-002"),
4654 ];
4655 let client_order_id1 = orders[0].client_order_id();
4656 let client_order_id2 = orders[1].client_order_id();
4657 let cache_rc = strategy.core.cache_rc();
4658 let timeline = Rc::new(RefCell::new(Vec::new()));
4659 let event_messages = Rc::new(RefCell::new(Vec::new()));
4660
4661 let event_handler = {
4662 let event_messages = event_messages.clone();
4663 let timeline = timeline.clone();
4664 TypedHandler::from_with_id("events.order.list_denied", move |event: &OrderEventAny| {
4665 match event {
4666 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id1 => {
4667 assert!(cache_rc.borrow().order_exists(&client_order_id1));
4668 timeline.borrow_mut().push("init1");
4669 }
4670 OrderEventAny::Initialized(e) if e.client_order_id == client_order_id2 => {
4671 assert!(cache_rc.borrow().order_exists(&client_order_id2));
4672 timeline.borrow_mut().push("init2");
4673 }
4674 OrderEventAny::Denied(e) if e.client_order_id == client_order_id1 => {
4675 assert_eq!(e.reason, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
4676 let cache = cache_rc.borrow();
4677 let cached_order = cache.order(&client_order_id1).unwrap();
4678 assert_eq!(cached_order.status(), OrderStatus::Denied);
4679 timeline.borrow_mut().push("denied1");
4680 }
4681 OrderEventAny::Denied(e) if e.client_order_id == client_order_id2 => {
4682 assert_eq!(e.reason, Ustr::from("MARKET_EXIT_IN_PROGRESS"));
4683 let cache = cache_rc.borrow();
4684 let cached_order = cache.order(&client_order_id2).unwrap();
4685 assert_eq!(cached_order.status(), OrderStatus::Denied);
4686 timeline.borrow_mut().push("denied2");
4687 }
4688 _ => panic!("unexpected order event {event:?}"),
4689 }
4690 event_messages.borrow_mut().push(event.clone());
4691 })
4692 };
4693 let risk_handler = {
4694 let timeline = timeline.clone();
4695 TypedIntoHandler::from_with_id(
4696 "RiskEngine.queue_execute",
4697 move |_command: TradingCommand| {
4698 timeline.borrow_mut().push("command");
4699 },
4700 )
4701 };
4702 msgbus::register_trading_command_endpoint(
4703 MessagingSwitchboard::risk_engine_queue_execute(),
4704 risk_handler,
4705 );
4706
4707 let topic = format!("events.order.{}", orders[0].strategy_id());
4708 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4709 let result = strategy.submit_order_list(orders.clone(), None, None, None);
4710
4711 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4712
4713 assert!(result.is_ok());
4714
4715 let cache = strategy.cache();
4716 let cached_order1 = cache.order(&client_order_id1).unwrap();
4717 let cached_order2 = cache.order(&client_order_id2).unwrap();
4718 assert_eq!(cached_order1.status(), OrderStatus::Denied);
4719 assert_eq!(cached_order2.status(), OrderStatus::Denied);
4720
4721 let event_messages = event_messages.borrow();
4722 assert_eq!(event_messages.len(), 4);
4723 assert_eq!(
4724 event_messages[0],
4725 OrderEventAny::Initialized(orders[0].init_event().clone())
4726 );
4727 assert!(matches!(
4728 &event_messages[1],
4729 OrderEventAny::Denied(e)
4730 if e.client_order_id == client_order_id1
4731 && e.reason == Ustr::from("MARKET_EXIT_IN_PROGRESS")
4732 ));
4733 assert_eq!(
4734 event_messages[2],
4735 OrderEventAny::Initialized(orders[1].init_event().clone())
4736 );
4737 assert!(matches!(
4738 &event_messages[3],
4739 OrderEventAny::Denied(e)
4740 if e.client_order_id == client_order_id2
4741 && e.reason == Ustr::from("MARKET_EXIT_IN_PROGRESS")
4742 ));
4743 assert_eq!(
4744 timeline.borrow().as_slice(),
4745 &["init1", "denied1", "init2", "denied2"]
4746 );
4747 }
4748
4749 #[rstest]
4750 fn test_submit_order_list_market_exit_rejects_non_initialized_without_events() {
4751 let mut strategy = create_test_strategy();
4752 register_strategy(&mut strategy);
4753 start_strategy(&mut strategy);
4754 strategy.core.is_exiting = true;
4755
4756 let order = make_accepted_market_order("O-20250208-LIST-DENY-ACCEPTED");
4757 let topic = format!("events.order.{}", order.strategy_id());
4758 let (event_handler, event_messages): (_, TypedMessageSavingHandler<OrderEventAny>) =
4759 get_typed_message_saving_handler(Some(Ustr::from("events.order.list_invalid")));
4760
4761 msgbus::subscribe_order_events(topic.clone().into(), event_handler.clone(), None);
4762 let result = strategy.submit_order_list(vec![order], None, None, None);
4763
4764 msgbus::unsubscribe_order_events(topic.into(), &event_handler);
4765
4766 assert!(result.is_err());
4767 assert!(
4768 result
4769 .unwrap_err()
4770 .to_string()
4771 .contains("expected INITIALIZED")
4772 );
4773 assert!(event_messages.get_messages().is_empty());
4774 }
4775
4776 #[rstest]
4777 fn test_submit_order_list_rejects_mixed_venues_with_friendly_error() {
4778 let mut strategy = create_test_strategy();
4779 register_strategy(&mut strategy);
4780 start_strategy(&mut strategy);
4781
4782 let binance_order = make_initialized_market_order("O-MIXED-VENUE-001");
4783 let bybit_order = OrderAny::Market(MarketOrder::new(
4784 TraderId::from("TRADER-001"),
4785 StrategyId::from("TEST-001"),
4786 InstrumentId::from("BTCUSDT.BYBIT"),
4787 ClientOrderId::from("O-MIXED-VENUE-002"),
4788 OrderSide::Buy,
4789 Quantity::from(100_000),
4790 TimeInForce::Gtc,
4791 UUID4::new(),
4792 UnixNanos::default(),
4793 false,
4794 false,
4795 None,
4796 None,
4797 None,
4798 None,
4799 None,
4800 None,
4801 None,
4802 None,
4803 ));
4804
4805 let result = strategy.submit_order_list(vec![binance_order, bybit_order], None, None, None);
4806
4807 let err = result.unwrap_err();
4808 let msg = err.to_string();
4809 assert!(
4810 msg.contains("OrderList denied: orders must share the same venue"),
4811 "unexpected error: {msg}",
4812 );
4813 assert!(msg.contains("BINANCE"), "expected BINANCE in error: {msg}");
4814 assert!(msg.contains("BYBIT"), "expected BYBIT in error: {msg}");
4815 }
4816
4817 #[rstest]
4818 fn test_submit_order_allowed_during_market_exit_when_reduce_only() {
4819 let mut strategy = create_test_strategy();
4820 register_strategy(&mut strategy);
4821 start_strategy(&mut strategy);
4822 strategy.core.is_exiting = true;
4823
4824 let order = OrderAny::Market(MarketOrder::new(
4825 TraderId::from("TRADER-001"),
4826 StrategyId::from("TEST-001"),
4827 InstrumentId::from("BTCUSDT.BINANCE"),
4828 ClientOrderId::from("O-20250208-0001"),
4829 OrderSide::Buy,
4830 Quantity::from(100_000),
4831 TimeInForce::Gtc,
4832 UUID4::new(),
4833 UnixNanos::default(),
4834 true, false,
4836 None,
4837 None,
4838 None,
4839 None,
4840 None,
4841 None,
4842 None,
4843 None,
4844 ));
4845 let client_order_id = order.client_order_id();
4846 let result = strategy.submit_order(order, None, None, None);
4847
4848 assert!(result.is_ok());
4849 let cache = strategy.cache();
4850 let cached_order = cache.order(&client_order_id).unwrap();
4851 assert_ne!(cached_order.status(), OrderStatus::Denied);
4852 }
4853
4854 #[rstest]
4855 fn test_submit_order_allowed_during_market_exit_when_tagged() {
4856 let mut strategy = create_test_strategy();
4857 register_strategy(&mut strategy);
4858 start_strategy(&mut strategy);
4859 strategy.core.is_exiting = true;
4860
4861 let order = OrderAny::Market(MarketOrder::new(
4862 TraderId::from("TRADER-001"),
4863 StrategyId::from("TEST-001"),
4864 InstrumentId::from("BTCUSDT.BINANCE"),
4865 ClientOrderId::from("O-20250208-0002"),
4866 OrderSide::Buy,
4867 Quantity::from(100_000),
4868 TimeInForce::Gtc,
4869 UUID4::new(),
4870 UnixNanos::default(),
4871 false, false,
4873 None,
4874 None,
4875 None,
4876 None,
4877 None,
4878 None,
4879 None,
4880 Some(vec![Ustr::from("MARKET_EXIT")]),
4881 ));
4882 let client_order_id = order.client_order_id();
4883 let result = strategy.submit_order(order, None, None, None);
4884
4885 assert!(result.is_ok());
4886 let cache = strategy.cache();
4887 let cached_order = cache.order(&client_order_id).unwrap();
4888 assert_ne!(cached_order.status(), OrderStatus::Denied);
4889 }
4890
4891 #[derive(Debug)]
4892 struct MacroTestSimple {
4893 core: StrategyCore,
4894 }
4895
4896 nautilus_strategy!(MacroTestSimple);
4897
4898 impl DataActor for MacroTestSimple {}
4899
4900 #[derive(Debug)]
4901 struct MacroTestWithHooks {
4902 core: StrategyCore,
4903 }
4904
4905 nautilus_strategy!(MacroTestWithHooks, {
4906 fn on_order_rejected(&mut self, _event: OrderRejected) {}
4907 });
4908
4909 impl DataActor for MacroTestWithHooks {}
4910
4911 #[derive(Debug)]
4912 struct MacroTestCustomField {
4913 inner: StrategyCore,
4914 }
4915
4916 nautilus_strategy!(MacroTestCustomField, inner, {
4917 fn external_order_claims(&self) -> Option<Vec<InstrumentId>> {
4918 None
4919 }
4920 });
4921
4922 impl DataActor for MacroTestCustomField {}
4923
4924 #[rstest]
4925 fn test_strategy_behavior_does_not_require_native_core_access() {
4926 fn assert_strategy<T: Strategy + DataActor + Component>() {}
4927
4928 assert_strategy::<CoreFreeStrategy>();
4929
4930 let mut strategy = CoreFreeStrategy {
4931 state: ComponentState::PreInitialized,
4932 started: false,
4933 };
4934 DataActor::on_start(&mut strategy).unwrap();
4935
4936 assert!(strategy.started);
4937 }
4938
4939 #[rstest]
4940 fn test_nautilus_strategy_macro_forms() {
4941 let config = StrategyConfig {
4942 strategy_id: Some(StrategyId::from("MACRO-001")),
4943 order_id_tag: Some("001".to_string()),
4944 ..Default::default()
4945 };
4946
4947 let simple = MacroTestSimple {
4948 core: StrategyCore::new(config.clone()),
4949 };
4950 assert_eq!(simple.strategy_id(), config.strategy_id);
4951 assert_eq!(simple.config().order_id_tag, config.order_id_tag);
4952 assert_eq!(simple.actor_id(), ActorId::from("MACRO-001"));
4953
4954 let hooks = MacroTestWithHooks {
4955 core: StrategyCore::new(config.clone()),
4956 };
4957 assert_eq!(hooks.strategy_id(), config.strategy_id);
4958 assert_eq!(hooks.config().order_id_tag, config.order_id_tag);
4959 assert_eq!(hooks.actor_id(), ActorId::from("MACRO-001"));
4960
4961 let custom = MacroTestCustomField {
4962 inner: StrategyCore::new(config.clone()),
4963 };
4964 assert_eq!(custom.strategy_id(), config.strategy_id);
4965 assert_eq!(custom.config().order_id_tag, config.order_id_tag);
4966 assert_eq!(custom.actor_id(), ActorId::from("MACRO-001"));
4967 assert!(custom.external_order_claims().is_none());
4968 }
4969}