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