1use std::time::Duration;
34
35use ahash::AHashMap;
36use nautilus_common::{
37 actor::{DataActor, DataActorNative},
38 timer::TimeEvent,
39};
40use nautilus_model::{
41 enums::OrderType,
42 events::OrderDeniedReason,
43 identifiers::ClientOrderId,
44 instruments::Instrument,
45 orders::{Order, OrderAny},
46 types::Quantity,
47};
48use rust_decimal::{Decimal, RoundingStrategy};
49use ustr::Ustr;
50
51use super::{
52 ExecutionAlgorithm, ExecutionAlgorithmConfig, ExecutionAlgorithmCore, ExecutionAlgorithmNative,
53};
54use crate::nautilus_execution_algorithm;
55
56pub type TwapAlgorithmConfig = ExecutionAlgorithmConfig;
58
59#[derive(Debug)]
65pub struct TwapAlgorithm {
66 pub core: ExecutionAlgorithmCore,
68 scheduled_orders: AHashMap<ClientOrderId, TwapSchedule>,
70}
71
72impl TwapAlgorithm {
73 #[must_use]
75 pub fn new(config: TwapAlgorithmConfig) -> Self {
76 Self {
77 core: ExecutionAlgorithmCore::new(config),
78 scheduled_orders: AHashMap::new(),
79 }
80 }
81
82 fn complete_sequence(&mut self, primary_id: ClientOrderId) {
84 let timer_name = primary_id.as_str();
85 let core = ExecutionAlgorithmNative::exec_algorithm_core_mut(self);
86 if core.clock_mut().timer_names().contains(&timer_name) {
87 core.clock_mut().cancel_timer(timer_name);
88 }
89 core.remove_submit_params(&primary_id);
90 self.scheduled_orders.remove(&primary_id);
91 log::info!("Completed TWAP execution for {primary_id}");
92 }
93}
94
95impl DataActor for TwapAlgorithm {
98 fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()> {
99 ExecutionAlgorithm::on_time_event(self, event)
100 }
101
102 fn on_stop(&mut self) -> anyhow::Result<()> {
103 ExecutionAlgorithm::on_stop(self)
104 }
105
106 fn on_resume(&mut self) -> anyhow::Result<()> {
107 ExecutionAlgorithm::on_resume(self)
108 }
109
110 fn on_reset(&mut self) -> anyhow::Result<()> {
111 ExecutionAlgorithm::on_reset(self)
112 }
113}
114
115nautilus_execution_algorithm!(TwapAlgorithm, {
116 fn on_order(&mut self, order: OrderAny) -> anyhow::Result<()> {
117 let primary_id = order.client_order_id();
118
119 if self.scheduled_orders.contains_key(&primary_id) {
120 anyhow::bail!("Order {primary_id} already being executed");
121 }
122
123 log::info!("Received order for TWAP execution: {order:?}");
124
125 if order.order_type() != OrderType::Market {
127 let reason = OrderDeniedReason::UnsupportedOrderType {
128 order_type: order.order_type(),
129 }
130 .to_string();
131 return self.deny_order(&order, Ustr::from(&reason));
132 }
133
134 let instrument = {
135 let cache = ExecutionAlgorithmNative::exec_algorithm_core(self).cache_ref();
136 cache.instrument(&order.instrument_id()).cloned()
137 };
138
139 let Some(instrument) = instrument else {
140 let reason = OrderDeniedReason::InstrumentNotFound {
141 instrument_id: order.instrument_id(),
142 }
143 .to_string();
144 return self.deny_order(&order, Ustr::from(&reason));
145 };
146
147 let Some(exec_params) = order.exec_algorithm_params() else {
148 return self.deny_order(&order, validation_failed("exec_algorithm_params not found"));
149 };
150
151 let Some(horizon_secs_str) = exec_params.get(&Ustr::from("horizon_secs")) else {
152 return self.deny_order(
153 &order,
154 validation_failed("horizon_secs not found in exec_algorithm_params"),
155 );
156 };
157
158 let horizon_secs: f64 = match horizon_secs_str.parse() {
159 Ok(value) => value,
160 Err(_) => {
161 return self.deny_order(
162 &order,
163 validation_failed(format!(
164 "horizon_secs={horizon_secs_str} is not a valid number"
165 )),
166 );
167 }
168 };
169
170 let Some(interval_secs_str) = exec_params.get(&Ustr::from("interval_secs")) else {
171 return self.deny_order(
172 &order,
173 validation_failed("interval_secs not found in exec_algorithm_params"),
174 );
175 };
176
177 let interval_secs: f64 = match interval_secs_str.parse() {
178 Ok(value) => value,
179 Err(_) => {
180 return self.deny_order(
181 &order,
182 validation_failed(format!(
183 "interval_secs={interval_secs_str} is not a valid number"
184 )),
185 );
186 }
187 };
188
189 if !horizon_secs.is_finite() || horizon_secs <= 0.0 {
190 return self.deny_order(
191 &order,
192 validation_failed(format!(
193 "horizon_secs={horizon_secs} must be finite and positive"
194 )),
195 );
196 }
197
198 if !interval_secs.is_finite() || interval_secs <= 0.0 {
199 return self.deny_order(
200 &order,
201 validation_failed(format!(
202 "interval_secs={interval_secs} must be finite and positive"
203 )),
204 );
205 }
206
207 if horizon_secs < interval_secs {
208 return self.deny_order(
209 &order,
210 validation_failed(format!(
211 "horizon_secs={horizon_secs} must be greater than or equal to interval_secs={interval_secs}"
212 )),
213 );
214 }
215
216 let num_intervals = (horizon_secs / interval_secs).floor() as u64;
217 if num_intervals == 0 {
218 return self.deny_order(&order, validation_failed("num_intervals is 0"));
219 }
220
221 let total_qty = order.quantity();
222 let interval_count = Decimal::from(num_intervals);
223 let quotient = total_qty.as_decimal() / interval_count;
224 let floored = quotient.round_dp_with_strategy(
225 u32::from(instrument.size_precision()),
226 RoundingStrategy::ToZero,
227 );
228 let qty_per_interval = match instrument.try_make_qty_from_decimal(floored, None) {
229 Ok(quantity) => quantity,
230 Err(e) => {
231 return self.deny_order(
232 &order,
233 validation_failed(format!("invalid qty_per_interval={floored}: {e}")),
234 );
235 }
236 };
237 let remainder = total_qty.as_decimal() - floored * interval_count;
238
239 if qty_per_interval == total_qty || qty_per_interval < instrument.size_increment() {
240 log::warn!(
241 "Submitting for entire size: qty_per_interval={qty_per_interval}, order_quantity={total_qty}"
242 );
243 self.submit_order(order, None, None)?;
244 self.complete_sequence(primary_id);
245 return Ok(());
246 }
247
248 if let Some(min_qty) = instrument.min_quantity()
249 && qty_per_interval < min_qty
250 {
251 log::warn!(
252 "Submitting for entire size: qty_per_interval={qty_per_interval} < min_quantity={min_qty}"
253 );
254 self.submit_order(order, None, None)?;
255 self.complete_sequence(primary_id);
256 return Ok(());
257 }
258
259 let interval = match Duration::try_from_secs_f64(interval_secs) {
260 Ok(interval) => interval,
261 Err(e) => {
262 return self.deny_order(
263 &order,
264 validation_failed(format!(
265 "interval_secs={interval_secs} is not a valid duration: {e}"
266 )),
267 );
268 }
269 };
270
271 if interval == Duration::ZERO {
272 return self.deny_order(
273 &order,
274 validation_failed(format!(
275 "interval_secs={interval_secs} rounds to a zero duration"
276 )),
277 );
278 }
279 let Ok(interval_ns) = u64::try_from(interval.as_nanos()) else {
280 return self.deny_order(
281 &order,
282 validation_failed(format!(
283 "interval_secs={interval_secs} exceeds the clock nanosecond range"
284 )),
285 );
286 };
287 let timestamp_ns = ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
288 .clock_mut()
289 .timestamp_ns()
290 .as_u64();
291
292 if timestamp_ns.checked_add(interval_ns).is_none() {
293 return self.deny_order(
294 &order,
295 validation_failed(format!(
296 "interval_secs={interval_secs} exceeds the clock timestamp headroom"
297 )),
298 );
299 }
300
301 let mut scheduled_sizes: Vec<Quantity> = vec![qty_per_interval; num_intervals as usize];
302
303 if remainder > Decimal::ZERO {
304 let remainder_qty = match instrument.try_make_qty_from_decimal(remainder, None) {
305 Ok(quantity) => quantity,
306 Err(e) => {
307 return self.deny_order(
308 &order,
309 validation_failed(format!("invalid qty_remainder={remainder}: {e}")),
310 );
311 }
312 };
313 scheduled_sizes.push(remainder_qty);
314 }
315
316 let scheduled_total = scheduled_sizes
317 .iter()
318 .fold(Decimal::ZERO, |total, quantity| {
319 total + quantity.as_decimal()
320 });
321
322 if scheduled_total != total_qty.as_decimal() {
323 return self.deny_order(
324 &order,
325 validation_failed(format!(
326 "scheduled quantity {scheduled_total} does not equal order quantity {total_qty}"
327 )),
328 );
329 }
330
331 log::info!("Order execution size schedule: {scheduled_sizes:?}");
332
333 {
336 let cache_rc = ExecutionAlgorithmNative::exec_algorithm_core(self).cache_rc();
337 let mut cache = cache_rc.borrow_mut();
338 if !cache.order_exists(&primary_id) {
339 cache.add_order(order.clone(), None, None, false)?;
340 }
341 }
342
343 self.scheduled_orders.insert(
344 primary_id,
345 TwapSchedule {
346 remaining_sizes: scheduled_sizes.clone(),
347 interval,
348 },
349 );
350
351 let schedule = self.scheduled_orders.get_mut(&primary_id).unwrap();
352 let first_qty = schedule.remaining_sizes.remove(0);
353 let is_single_slice = self
354 .scheduled_orders
355 .get(&primary_id)
356 .is_some_and(|schedule| schedule.remaining_sizes.is_empty());
357
358 if is_single_slice {
360 self.submit_order(order, None, None)?;
361 self.complete_sequence(primary_id);
362 return Ok(());
363 }
364
365 let tags = order.tags().map(<[Ustr]>::to_vec);
367 let time_in_force = order.time_in_force();
368 let reduce_only = order.is_reduce_only();
369 let mut order = order;
370 let spawned = self.spawn_market(
371 &mut order,
372 first_qty,
373 time_in_force,
374 reduce_only,
375 tags,
376 true,
377 );
378 self.submit_order(spawned.into(), None, None)?;
379
380 ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
381 .clock_mut()
382 .set_timer(primary_id.as_str(), interval, None, None, None, None, None)?;
383
384 log::info!(
385 "Started TWAP execution for {primary_id}: horizon_secs={horizon_secs}, interval_secs={interval_secs}"
386 );
387
388 Ok(())
389 }
390
391 fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()> {
392 log::info!("Received time event: {event:?}");
393
394 let primary_id = ClientOrderId::new(event.name.as_str());
395
396 let primary = {
397 let cache = ExecutionAlgorithmNative::exec_algorithm_core(self).cache_ref();
398 cache.order(&primary_id).map(|o| o.clone())
399 };
400
401 let Some(primary) = primary else {
402 log::error!("Cannot find primary order for exec_spawn_id={primary_id}");
403 self.complete_sequence(primary_id);
404 return Ok(());
405 };
406
407 if primary.is_closed() {
408 self.complete_sequence(primary_id);
409 return Ok(());
410 }
411
412 let Some(schedule) = self.scheduled_orders.get_mut(&primary_id) else {
413 log::error!("Cannot find scheduled sizes for exec_spawn_id={primary_id}");
414 return Ok(());
415 };
416
417 if schedule.remaining_sizes.is_empty() {
418 log::warn!("No more size to execute for exec_spawn_id={primary_id}");
419 return Ok(());
420 }
421
422 let quantity = schedule.remaining_sizes.remove(0);
423 let is_final_slice = schedule.remaining_sizes.is_empty();
424
425 if is_final_slice {
427 self.submit_order(primary, None, None)?;
428 self.complete_sequence(primary_id);
429 return Ok(());
430 }
431
432 let tags = primary.tags().map(<[Ustr]>::to_vec);
434 let time_in_force = primary.time_in_force();
435 let reduce_only = primary.is_reduce_only();
436 let mut primary = primary;
437 let spawned = self.spawn_market(
438 &mut primary,
439 quantity,
440 time_in_force,
441 reduce_only,
442 tags,
443 true,
444 );
445 self.submit_order(spawned.into(), None, None)?;
446
447 Ok(())
448 }
449
450 fn on_stop(&mut self) -> anyhow::Result<()> {
451 ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
452 .clock_mut()
453 .cancel_timers();
454 Ok(())
455 }
456
457 fn on_resume(&mut self) -> anyhow::Result<()> {
458 let primary_ids: Vec<ClientOrderId> = self.scheduled_orders.keys().copied().collect();
459
460 for primary_id in primary_ids {
461 let primary_is_open = {
462 let cache = ExecutionAlgorithmNative::exec_algorithm_core(self).cache_ref();
463 cache.order(&primary_id).map(|primary| !primary.is_closed())
464 };
465
466 if primary_is_open.is_none() {
467 log::error!("Cannot find primary order for exec_spawn_id={primary_id}");
468 }
469 let interval = self.scheduled_orders.get(&primary_id).and_then(|schedule| {
470 (!schedule.remaining_sizes.is_empty()).then_some(schedule.interval)
471 });
472
473 let Some(interval) = interval.filter(|_| primary_is_open == Some(true)) else {
474 self.complete_sequence(primary_id);
475 continue;
476 };
477
478 ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
479 .clock_mut()
480 .set_timer(primary_id.as_str(), interval, None, None, None, None, None)?;
481 }
482
483 Ok(())
484 }
485
486 fn on_reset(&mut self) -> anyhow::Result<()> {
487 self.unsubscribe_all_strategy_events();
488 ExecutionAlgorithmNative::exec_algorithm_core_mut(self).reset();
489 self.scheduled_orders.clear();
490 Ok(())
491 }
492});
493
494#[derive(Debug)]
495struct TwapSchedule {
496 remaining_sizes: Vec<Quantity>,
497 interval: Duration,
498}
499
500fn validation_failed(detail: impl Into<String>) -> Ustr {
501 let reason = OrderDeniedReason::ValidationFailed {
502 detail: detail.into(),
503 }
504 .to_string();
505 Ustr::from(&reason)
506}
507
508#[cfg(test)]
509mod tests {
510 use std::{cell::RefCell, rc::Rc};
511
512 use indexmap::IndexMap;
513 use nautilus_common::{
514 cache::Cache,
515 clock::{Clock, TestClock},
516 component::Component,
517 enums::ComponentTrigger,
518 messages::execution::{SubmitOrder, TradingCommand},
519 msgbus::{self, MessagingSwitchboard, TypedHandler},
520 };
521 use nautilus_core::{Params, UUID4, UnixNanos};
522 use nautilus_model::{
523 enums::{OrderSide, OrderStatus, TimeInForce},
524 events::{OrderDeniedReason, OrderEventAny, order::spec::OrderCanceledSpec},
525 identifiers::{ExecAlgorithmId, InstrumentId, StrategyId, TraderId},
526 orders::{LimitOrder, MarketOrder},
527 types::Price,
528 };
529 use rstest::rstest;
530 use ustr::Ustr;
531
532 use super::*;
533
534 fn create_twap_algorithm() -> TwapAlgorithm {
535 let unique_id = format!("TWAP-{}", UUID4::new());
537 let config = TwapAlgorithmConfig {
538 exec_algorithm_id: Some(ExecAlgorithmId::new(&unique_id)),
539 ..Default::default()
540 };
541 TwapAlgorithm::new(config)
542 }
543
544 fn register_algorithm_with_clock(algo: &mut TwapAlgorithm) -> Rc<RefCell<TestClock>> {
545 use nautilus_common::timer::TimeEventCallback;
546
547 let trader_id = TraderId::from("TRADER-001");
548 let clock = Rc::new(RefCell::new(TestClock::new()));
549 let cache = Rc::new(RefCell::new(Cache::default()));
550
551 clock
553 .borrow_mut()
554 .register_default_handler(TimeEventCallback::Rust(std::sync::Arc::new(|_| {})));
555
556 algo.core.register(trader_id, clock.clone(), cache).unwrap();
557
558 algo.transition_state(ComponentTrigger::Initialize).unwrap();
560 algo.transition_state(ComponentTrigger::Start).unwrap();
561 algo.transition_state(ComponentTrigger::StartCompleted)
562 .unwrap();
563
564 clock
565 }
566
567 fn register_algorithm(algo: &mut TwapAlgorithm) {
568 let _ = register_algorithm_with_clock(algo);
569 }
570
571 fn add_instrument_to_cache(algo: &TwapAlgorithm) {
572 use nautilus_model::instruments::{InstrumentAny, stubs::crypto_perpetual_ethusdt};
573
574 let instrument = crypto_perpetual_ethusdt();
575 let cache_rc = algo.core.cache_rc();
576 let mut cache = cache_rc.borrow_mut();
577 cache
578 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
579 .unwrap();
580 }
581
582 fn create_market_order_with_params(params: IndexMap<Ustr, Ustr>) -> OrderAny {
583 create_market_order_with_params_and_qty(params, Quantity::from("1.0"))
584 }
585
586 fn create_market_order_with_params_and_qty(
587 params: IndexMap<Ustr, Ustr>,
588 quantity: Quantity,
589 ) -> OrderAny {
590 let client_order_id = ClientOrderId::from("O-001");
591 OrderAny::Market(MarketOrder::new(
592 TraderId::from("TRADER-001"),
593 StrategyId::from("STRAT-001"),
594 InstrumentId::from("ETHUSDT-PERP.BINANCE"),
595 client_order_id,
596 OrderSide::Buy,
597 quantity,
598 TimeInForce::Gtc,
599 UUID4::new(),
600 0.into(),
601 false,
602 false,
603 None,
604 None,
605 None,
606 None,
607 Some(ExecAlgorithmId::new("TWAP")),
608 Some(params),
609 Some(client_order_id),
610 None,
611 ))
612 }
613
614 fn assert_twap_denied(algo: &mut TwapAlgorithm, order: &OrderAny, expected_reason: &str) {
615 let strategy_id = order.strategy_id();
616 {
617 let cache_rc = algo.core.cache_rc();
618 cache_rc
619 .borrow_mut()
620 .add_order(order.clone(), None, None, false)
621 .unwrap();
622 }
623 let events = Rc::new(RefCell::new(Vec::new()));
624 let handler = TypedHandler::from({
625 let events = events.clone();
626 move |event: &OrderEventAny| events.borrow_mut().push(event.clone())
627 });
628 let topic = format!("events.order.{strategy_id}");
629 msgbus::subscribe_order_events(topic.clone().into(), handler.clone(), None);
630
631 algo.on_order(order.clone()).unwrap();
632 algo.on_order(order.clone()).unwrap();
633
634 msgbus::unsubscribe_order_events(topic.into(), &handler);
635 let cached_order = algo.cache().order(&order.client_order_id()).unwrap();
636 let events = events.borrow();
637
638 assert_eq!(cached_order.status(), OrderStatus::Denied);
639 assert_eq!(events.len(), 1);
640 assert!(matches!(
641 &events[0],
642 OrderEventAny::Denied(event)
643 if event.reason.as_str() == expected_reason
644 && event.strategy_id == strategy_id
645 && event.client_order_id == order.client_order_id()
646 ));
647 assert!(algo.scheduled_orders.is_empty());
648 assert!(algo.clock().timer_names().is_empty());
649 }
650
651 #[rstest]
652 fn test_twap_creation() {
653 let algo = create_twap_algorithm();
654 assert!(algo.id().inner().starts_with("TWAP"));
655 assert!(algo.scheduled_orders.is_empty());
656 }
657
658 #[rstest]
659 fn test_twap_registration() {
660 let mut algo = create_twap_algorithm();
661 register_algorithm(&mut algo);
662
663 assert_eq!(algo.trader_id(), Some(TraderId::from("TRADER-001")));
664 }
665
666 #[rstest]
667 fn test_twap_reset_clears_scheduled_sizes() {
668 let mut algo = create_twap_algorithm();
669 algo.scheduled_orders.insert(
670 ClientOrderId::new("O-001"),
671 TwapSchedule {
672 remaining_sizes: vec![Quantity::from("1.0")],
673 interval: Duration::from_secs(1),
674 },
675 );
676 algo.scheduled_orders.insert(
677 ClientOrderId::new("O-002"),
678 TwapSchedule {
679 remaining_sizes: vec![Quantity::from("2.0")],
680 interval: Duration::from_secs(2),
681 },
682 );
683
684 assert!(!algo.scheduled_orders.is_empty());
685
686 DataActor::on_reset(&mut algo).unwrap();
688
689 assert!(algo.scheduled_orders.is_empty());
690 }
691
692 #[rstest]
693 fn test_twap_rejects_non_market_orders() {
694 let mut algo = create_twap_algorithm();
695 register_algorithm(&mut algo);
696
697 let order = OrderAny::Limit(LimitOrder::new(
698 TraderId::from("TRADER-001"),
699 StrategyId::from("STRAT-001"),
700 InstrumentId::from("BTC/USDT.BINANCE"),
701 ClientOrderId::from("O-001"),
702 OrderSide::Buy,
703 Quantity::from("1.0"),
704 Price::from("50000.0"),
705 TimeInForce::Gtc,
706 None, false, false, false, None, None, None, None, None, None, None, None, None, None, None, UUID4::new(),
722 0.into(),
723 ));
724
725 let reason = OrderDeniedReason::UnsupportedOrderType {
726 order_type: OrderType::Limit,
727 }
728 .to_string();
729 assert_twap_denied(&mut algo, &order, &reason);
730 }
731
732 #[rstest]
733 fn test_twap_denies_missing_instrument() {
734 let mut algo = create_twap_algorithm();
735 register_algorithm(&mut algo);
736
737 let mut params = IndexMap::new();
738 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
739 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
740 let order = create_market_order_with_params(params);
741 let reason = OrderDeniedReason::InstrumentNotFound {
742 instrument_id: order.instrument_id(),
743 }
744 .to_string();
745
746 assert_twap_denied(&mut algo, &order, &reason);
747 }
748
749 #[rstest]
750 fn test_twap_rejects_missing_params() {
751 let mut algo = create_twap_algorithm();
752 register_algorithm(&mut algo);
753
754 add_instrument_to_cache(&algo);
755
756 let order = OrderAny::Market(MarketOrder::new(
757 TraderId::from("TRADER-001"),
758 StrategyId::from("STRAT-001"),
759 InstrumentId::from("ETHUSDT-PERP.BINANCE"),
760 ClientOrderId::from("O-001"),
761 OrderSide::Buy,
762 Quantity::from("1.0"),
763 TimeInForce::Gtc,
764 UUID4::new(),
765 0.into(),
766 false,
767 false,
768 None,
769 None,
770 None,
771 None,
772 None,
773 None, None,
775 None,
776 ));
777
778 assert_twap_denied(
779 &mut algo,
780 &order,
781 "VALIDATION_FAILED: exec_algorithm_params not found",
782 );
783 }
784
785 #[rstest]
786 #[case(
787 None,
788 Some("10"),
789 "VALIDATION_FAILED: horizon_secs not found in exec_algorithm_params"
790 )]
791 #[case(
792 Some("60"),
793 None,
794 "VALIDATION_FAILED: interval_secs not found in exec_algorithm_params"
795 )]
796 #[case(
797 Some("not-a-number"),
798 Some("10"),
799 "VALIDATION_FAILED: horizon_secs=not-a-number is not a valid number"
800 )]
801 #[case(
802 Some("60"),
803 Some("not-a-number"),
804 "VALIDATION_FAILED: interval_secs=not-a-number is not a valid number"
805 )]
806 fn test_twap_denies_missing_or_malformed_schedule_parameter(
807 #[case] horizon_secs: Option<&str>,
808 #[case] interval_secs: Option<&str>,
809 #[case] expected_reason: &str,
810 ) {
811 let mut algo = create_twap_algorithm();
812 register_algorithm(&mut algo);
813 add_instrument_to_cache(&algo);
814
815 let mut params = IndexMap::new();
816
817 if let Some(horizon_secs) = horizon_secs {
818 params.insert(Ustr::from("horizon_secs"), Ustr::from(horizon_secs));
819 }
820
821 if let Some(interval_secs) = interval_secs {
822 params.insert(Ustr::from("interval_secs"), Ustr::from(interval_secs));
823 }
824
825 assert_twap_denied(
826 &mut algo,
827 &create_market_order_with_params(params),
828 expected_reason,
829 );
830 }
831
832 #[rstest]
833 fn test_twap_rejects_horizon_less_than_interval() {
834 let mut algo = create_twap_algorithm();
835 register_algorithm(&mut algo);
836
837 add_instrument_to_cache(&algo);
838
839 let mut params = IndexMap::new();
840 params.insert(Ustr::from("horizon_secs"), Ustr::from("30"));
841 params.insert(Ustr::from("interval_secs"), Ustr::from("60"));
842
843 let order = create_market_order_with_params(params);
844 assert_twap_denied(
845 &mut algo,
846 &order,
847 "VALIDATION_FAILED: horizon_secs=30 must be greater than or equal to interval_secs=60",
848 );
849 }
850
851 #[rstest]
852 fn test_twap_rejects_duplicate_order() {
853 let mut algo = create_twap_algorithm();
854 register_algorithm(&mut algo);
855
856 add_instrument_to_cache(&algo);
857
858 let mut params = IndexMap::new();
859 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
860 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
861
862 let order1 = create_market_order_with_params(params.clone());
863 let order2 = create_market_order_with_params(params);
864
865 algo.on_order(order1).unwrap();
866 let result = algo.on_order(order2);
867
868 assert!(result.is_err());
869 assert!(
870 result
871 .unwrap_err()
872 .to_string()
873 .contains("already being executed")
874 );
875 }
876
877 #[rstest]
878 fn test_twap_calculates_size_schedule_evenly() {
879 let mut algo = create_twap_algorithm();
880 register_algorithm(&mut algo);
881
882 add_instrument_to_cache(&algo);
883
884 let mut params = IndexMap::new();
886 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
887 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
888
889 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
890 let primary_id = order.client_order_id();
891
892 algo.on_order(order).unwrap();
893
894 let remaining = &algo
896 .scheduled_orders
897 .get(&primary_id)
898 .unwrap()
899 .remaining_sizes;
900 assert_eq!(remaining.len(), 2);
901
902 for qty in remaining {
903 assert_eq!(*qty, Quantity::from("0.4"));
904 }
905 }
906
907 #[rstest]
908 fn test_twap_reduces_cached_primary_after_first_child_spawn() {
909 let mut algo = create_twap_algorithm();
910 register_algorithm(&mut algo);
911
912 add_instrument_to_cache(&algo);
913
914 let mut params = IndexMap::new();
915 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
916 params.insert(Ustr::from("interval_secs"), Ustr::from("30"));
917
918 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
919 let primary_id = order.client_order_id();
920
921 algo.on_order(order).unwrap();
922
923 let cache = algo.cache();
924 let primary = cache.order(&primary_id).unwrap();
925 let spawned = cache.order(&ClientOrderId::from("O-001-E1")).unwrap();
926
927 assert_eq!(primary.quantity(), Quantity::from("0.6"));
928 assert_eq!(spawned.quantity(), Quantity::from("0.6"));
929 assert_eq!(spawned.exec_spawn_id(), Some(primary_id));
930 }
931
932 #[rstest]
933 fn test_twap_on_order_accepts_already_cached_primary() {
934 let mut algo = create_twap_algorithm();
935 register_algorithm(&mut algo);
936
937 add_instrument_to_cache(&algo);
938
939 let mut params = IndexMap::new();
940 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
941 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
942
943 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
944 let primary_id = order.client_order_id();
945
946 {
948 let cache_rc = algo.core.cache_rc();
949 let mut cache = cache_rc.borrow_mut();
950 cache.add_order(order.clone(), None, None, false).unwrap();
951 }
952
953 algo.on_order(order).unwrap();
954
955 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 2);
956 }
957
958 #[rstest]
959 fn test_twap_calculates_size_schedule_with_remainder() {
960 let mut algo = create_twap_algorithm();
961 register_algorithm(&mut algo);
962
963 add_instrument_to_cache(&algo);
964 let instrument = nautilus_model::instruments::stubs::crypto_perpetual_ethusdt();
965
966 let mut params = IndexMap::new();
968 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
969 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
970
971 let order = create_market_order_with_params(params);
972 let primary_id = order.client_order_id();
973
974 algo.on_order(order).unwrap();
975
976 let remaining = &algo
978 .scheduled_orders
979 .get(&primary_id)
980 .unwrap()
981 .remaining_sizes;
982 assert_eq!(remaining.len(), 3);
983 assert_eq!(
984 remaining,
985 &[
986 Quantity::from("0.333"),
987 Quantity::from("0.333"),
988 Quantity::from("0.001"),
989 ]
990 );
991
992 for quantity in remaining {
993 assert_eq!(quantity.precision, instrument.size_precision());
994 assert_eq!(quantity.raw % instrument.size_increment().raw, 0);
995 }
996
997 let first = algo
998 .cache()
999 .order(&ClientOrderId::from("O-001-E1"))
1000 .unwrap()
1001 .quantity();
1002 let total = remaining
1003 .iter()
1004 .fold(first.as_decimal(), |sum, qty| sum + qty.as_decimal());
1005 assert_eq!(total, Quantity::from("1.0").as_decimal());
1006 }
1007
1008 #[rstest]
1009 fn test_twap_children_use_instrument_size_precision() {
1010 let mut algo = create_twap_algorithm();
1011 register_algorithm(&mut algo);
1012
1013 add_instrument_to_cache(&algo);
1014 let instrument = nautilus_model::instruments::stubs::crypto_perpetual_ethusdt();
1015
1016 let mut params = IndexMap::new();
1017 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1018 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1019
1020 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.000000000"));
1021 let primary_id = order.client_order_id();
1022
1023 algo.on_order(order).unwrap();
1024
1025 let spawned = algo
1026 .cache()
1027 .order(&ClientOrderId::from("O-001-E1"))
1028 .unwrap()
1029 .quantity();
1030 assert_eq!(spawned.precision, instrument.size_precision());
1031 assert!(
1032 algo.scheduled_orders[&primary_id]
1033 .remaining_sizes
1034 .iter()
1035 .all(|quantity| quantity.precision == instrument.size_precision())
1036 );
1037 }
1038
1039 #[rstest]
1040 fn test_twap_on_time_event_spawns_next_slice() {
1041 let mut algo = create_twap_algorithm();
1042 register_algorithm(&mut algo);
1043
1044 add_instrument_to_cache(&algo);
1045
1046 let mut params = IndexMap::new();
1048 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1049 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1050
1051 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1052 let primary_id = order.client_order_id();
1053 let mut submit_params = Params::new();
1054 submit_params.insert(
1055 "routing_profile".to_string(),
1056 serde_json::Value::String("intermediate-slice".to_string()),
1057 );
1058 algo.core
1059 .remember_submit_params(primary_id, Some(submit_params.clone()));
1060
1061 algo.on_order(order).unwrap();
1062
1063 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 2);
1065
1066 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1068 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1069
1070 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 1);
1072 assert_eq!(algo.core.submit_params(&primary_id), Some(submit_params));
1073 }
1074
1075 #[rstest]
1076 fn test_twap_data_actor_dispatch_spawns_next_slice() {
1077 let mut algo = create_twap_algorithm();
1078 register_algorithm(&mut algo);
1079
1080 add_instrument_to_cache(&algo);
1081
1082 let mut params = IndexMap::new();
1083 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1084 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1085
1086 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1087 let primary_id = order.client_order_id();
1088
1089 algo.on_order(order).unwrap();
1090 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 2);
1091
1092 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1094 algo.handle_time_event(&event);
1095
1096 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 1);
1097 }
1098
1099 #[rstest]
1100 fn test_twap_on_time_event_completes_on_final_slice() {
1101 let mut algo = create_twap_algorithm();
1102 register_algorithm(&mut algo);
1103
1104 add_instrument_to_cache(&algo);
1105
1106 let mut params = IndexMap::new();
1108 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1109 params.insert(Ustr::from("interval_secs"), Ustr::from("30"));
1110
1111 let order = create_market_order_with_params(params);
1112 let primary_id = order.client_order_id();
1113 let mut submit_params = Params::new();
1114 submit_params.insert(
1115 "routing_profile".to_string(),
1116 serde_json::Value::String("final-slice".to_string()),
1117 );
1118 algo.core
1119 .remember_submit_params(primary_id, Some(submit_params.clone()));
1120
1121 algo.on_order(order).unwrap();
1122 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 1);
1123 assert_eq!(
1124 algo.core.submit_params(&primary_id),
1125 Some(submit_params.clone())
1126 );
1127
1128 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1129 let handler = msgbus::TypedIntoHandler::from({
1130 let captured = received.clone();
1131 move |cmd: TradingCommand| {
1132 if let TradingCommand::SubmitOrder(cmd) = cmd {
1133 *captured.borrow_mut() = Some(cmd);
1134 }
1135 }
1136 });
1137 msgbus::register_trading_command_endpoint(
1138 MessagingSwitchboard::risk_engine_queue_execute(),
1139 handler,
1140 );
1141
1142 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1144 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1145
1146 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1148 assert_eq!(
1149 received
1150 .borrow()
1151 .as_ref()
1152 .and_then(|cmd| cmd.params.clone()),
1153 Some(submit_params),
1154 );
1155 assert_eq!(algo.core.submit_params(&primary_id), None);
1156 }
1157
1158 #[rstest]
1159 fn test_twap_on_time_event_completes_when_primary_closed() {
1160 let mut algo = create_twap_algorithm();
1161 register_algorithm(&mut algo);
1162
1163 add_instrument_to_cache(&algo);
1164
1165 let mut params = IndexMap::new();
1166 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1167 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1168
1169 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1170 let primary_id = order.client_order_id();
1171 let mut submit_params = Params::new();
1172 submit_params.insert(
1173 "routing_profile".to_string(),
1174 serde_json::Value::String("closed-primary".to_string()),
1175 );
1176 algo.core
1177 .remember_submit_params(primary_id, Some(submit_params));
1178
1179 algo.on_order(order).unwrap();
1180 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 2);
1181
1182 {
1184 let cache_rc = algo.core.cache_rc();
1185 let mut cache = cache_rc.borrow_mut();
1186 let primary = cache.order(&primary_id).map(|o| o.clone()).unwrap();
1187
1188 let canceled = OrderCanceledSpec::builder()
1189 .trader_id(primary.trader_id())
1190 .strategy_id(primary.strategy_id())
1191 .instrument_id(primary.instrument_id())
1192 .client_order_id(primary.client_order_id())
1193 .build();
1194 cache
1195 .update_order(&OrderEventAny::Canceled(canceled))
1196 .unwrap();
1197 }
1198
1199 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1201 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1202
1203 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1205 assert_eq!(algo.core.submit_params(&primary_id), None);
1206 }
1207
1208 #[rstest]
1209 fn test_twap_on_time_event_completes_when_primary_missing() {
1210 let mut algo = create_twap_algorithm();
1211 register_algorithm(&mut algo);
1212 add_instrument_to_cache(&algo);
1213
1214 let mut params = IndexMap::new();
1215 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1216 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1217
1218 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1219 let primary_id = order.client_order_id();
1220 algo.on_order(order).unwrap();
1221
1222 {
1223 let cache_rc = algo.core.cache_rc();
1224 let mut cache = cache_rc.borrow_mut();
1225 let primary = cache.order(&primary_id).map(|order| order.clone()).unwrap();
1226 let canceled = OrderCanceledSpec::builder()
1227 .trader_id(primary.trader_id())
1228 .strategy_id(primary.strategy_id())
1229 .instrument_id(primary.instrument_id())
1230 .client_order_id(primary.client_order_id())
1231 .build();
1232 cache
1233 .update_order(&OrderEventAny::Canceled(canceled))
1234 .unwrap();
1235 cache.purge_order(primary_id);
1236 }
1237 assert!(algo.cache().order(&primary_id).is_none());
1238
1239 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1240 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1241
1242 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1244 assert!(algo.clock().timer_names().is_empty());
1245 }
1246
1247 #[rstest]
1248 fn test_twap_on_stop_cancels_timers() {
1249 let mut algo = create_twap_algorithm();
1250 register_algorithm(&mut algo);
1251
1252 add_instrument_to_cache(&algo);
1253
1254 let mut params = IndexMap::new();
1255 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1256 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1257
1258 let order = create_market_order_with_params(params);
1259 let primary_id = order.client_order_id();
1260
1261 algo.on_order(order).unwrap();
1262
1263 assert!(
1265 algo.clock()
1266 .timer_names()
1267 .iter()
1268 .any(|name| name.as_str() == primary_id.as_str())
1269 );
1270
1271 DataActor::on_stop(&mut algo).unwrap();
1273
1274 assert!(algo.clock().timer_names().is_empty());
1276 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 3);
1277 }
1278
1279 #[rstest]
1280 fn test_twap_on_resume_rearms_timer_without_submitting() {
1281 let mut algo = create_twap_algorithm();
1282 register_algorithm(&mut algo);
1283 add_instrument_to_cache(&algo);
1284
1285 let mut params = IndexMap::new();
1286 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1287 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1288
1289 let order = create_market_order_with_params(params);
1290 let primary_id = order.client_order_id();
1291 algo.on_order(order).unwrap();
1292 let order_count = algo
1293 .cache()
1294 .orders_total_count(None, None, None, None, None);
1295
1296 Component::stop(&mut algo).unwrap();
1297
1298 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1301 let handler = msgbus::TypedIntoHandler::from({
1302 let captured = received.clone();
1303 move |cmd: TradingCommand| {
1304 if let TradingCommand::SubmitOrder(cmd) = cmd {
1305 *captured.borrow_mut() = Some(cmd);
1306 }
1307 }
1308 });
1309 msgbus::register_trading_command_endpoint(
1310 MessagingSwitchboard::risk_engine_queue_execute(),
1311 handler,
1312 );
1313
1314 let resume_time = algo.clock().timestamp_ns();
1315 Component::resume(&mut algo).unwrap();
1316
1317 assert!(received.borrow().is_none());
1318 assert_eq!(algo.clock().timer_count(), 1);
1319 assert_eq!(
1320 algo.clock().next_time_ns(primary_id.as_str()),
1321 Some(resume_time + 20_000_000_000)
1322 );
1323 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 3);
1324 assert_eq!(
1325 algo.cache()
1326 .orders_total_count(None, None, None, None, None),
1327 order_count
1328 );
1329 }
1330
1331 #[rstest]
1332 fn test_twap_on_resume_executes_remaining_slices() {
1333 let mut algo = create_twap_algorithm();
1334 let clock = register_algorithm_with_clock(&mut algo);
1335 add_instrument_to_cache(&algo);
1336
1337 let mut params = IndexMap::new();
1338 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1339 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1340
1341 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1342 let primary_id = order.client_order_id();
1343 algo.on_order(order).unwrap();
1344
1345 Component::stop(&mut algo).unwrap();
1346 Component::resume(&mut algo).unwrap();
1347 assert_eq!(algo.clock().timer_count(), 1);
1348
1349 let first_events = clock.borrow_mut().advance_time(20_000_000_000.into(), true);
1350 assert_eq!(first_events.len(), 1);
1351 algo.handle_time_event(&first_events[0]);
1352 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 1);
1353
1354 let final_events = clock.borrow_mut().advance_time(40_000_000_000.into(), true);
1355 assert_eq!(final_events.len(), 1);
1356 algo.handle_time_event(&final_events[0]);
1357 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1358 }
1359
1360 #[rstest]
1361 fn test_twap_on_resume_completes_closed_primary() {
1362 let mut algo = create_twap_algorithm();
1363 register_algorithm(&mut algo);
1364 add_instrument_to_cache(&algo);
1365
1366 let mut params = IndexMap::new();
1367 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1368 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1369
1370 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1371 let primary_id = order.client_order_id();
1372 algo.on_order(order).unwrap();
1373 Component::stop(&mut algo).unwrap();
1374
1375 {
1376 let cache_rc = algo.core.cache_rc();
1377 let mut cache = cache_rc.borrow_mut();
1378 let primary = cache.order(&primary_id).map(|order| order.clone()).unwrap();
1379 let canceled = OrderCanceledSpec::builder()
1380 .trader_id(primary.trader_id())
1381 .strategy_id(primary.strategy_id())
1382 .instrument_id(primary.instrument_id())
1383 .client_order_id(primary.client_order_id())
1384 .build();
1385 cache
1386 .update_order(&OrderEventAny::Canceled(canceled))
1387 .unwrap();
1388 }
1389
1390 Component::resume(&mut algo).unwrap();
1391
1392 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1393 assert_eq!(algo.clock().timer_count(), 0);
1394 }
1395
1396 #[rstest]
1397 fn test_twap_on_resume_completes_missing_primary() {
1398 let mut algo = create_twap_algorithm();
1399 register_algorithm(&mut algo);
1400 add_instrument_to_cache(&algo);
1401
1402 let mut params = IndexMap::new();
1403 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1404 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1405
1406 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1407 let primary_id = order.client_order_id();
1408 algo.on_order(order).unwrap();
1409 Component::stop(&mut algo).unwrap();
1410
1411 {
1412 let cache_rc = algo.core.cache_rc();
1413 let mut cache = cache_rc.borrow_mut();
1414 let primary = cache.order(&primary_id).map(|order| order.clone()).unwrap();
1415 let canceled = OrderCanceledSpec::builder()
1416 .trader_id(primary.trader_id())
1417 .strategy_id(primary.strategy_id())
1418 .instrument_id(primary.instrument_id())
1419 .client_order_id(primary.client_order_id())
1420 .build();
1421 cache
1422 .update_order(&OrderEventAny::Canceled(canceled))
1423 .unwrap();
1424 cache.purge_order(primary_id);
1425 }
1426
1427 assert!(algo.cache().order(&primary_id).is_none());
1428 Component::resume(&mut algo).unwrap();
1429
1430 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1431 assert_eq!(algo.clock().timer_count(), 0);
1432 }
1433
1434 #[rstest]
1435 fn test_twap_fractional_interval_secs() {
1436 let mut algo = create_twap_algorithm();
1437 register_algorithm(&mut algo);
1438
1439 add_instrument_to_cache(&algo);
1440
1441 let mut params = IndexMap::new();
1443 params.insert(Ustr::from("horizon_secs"), Ustr::from("3"));
1444 params.insert(Ustr::from("interval_secs"), Ustr::from("0.5"));
1445
1446 let order = create_market_order_with_params(params);
1447 let primary_id = order.client_order_id();
1448
1449 algo.on_order(order).unwrap();
1451
1452 let remaining = &algo
1454 .scheduled_orders
1455 .get(&primary_id)
1456 .unwrap()
1457 .remaining_sizes;
1458 assert!(remaining.len() >= 5);
1459 }
1460
1461 #[rstest]
1462 fn test_twap_submits_entire_size_when_qty_per_interval_below_size_increment() {
1463 use nautilus_model::instruments::{InstrumentAny, stubs::equity_aapl};
1464
1465 let mut algo = create_twap_algorithm();
1466 register_algorithm(&mut algo);
1467
1468 let instrument = equity_aapl();
1470 let instrument_id = instrument.id();
1471 {
1472 let cache_rc = algo.core.cache_rc();
1473 let mut cache = cache_rc.borrow_mut();
1474 cache
1475 .add_instrument(InstrumentAny::Equity(instrument))
1476 .unwrap();
1477 }
1478
1479 let mut params = IndexMap::new();
1482 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1483 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
1484
1485 let client_order_id = ClientOrderId::from("O-002");
1486 let order = OrderAny::Market(MarketOrder::new(
1487 TraderId::from("TRADER-001"),
1488 StrategyId::from("STRAT-001"),
1489 instrument_id,
1490 client_order_id,
1491 OrderSide::Buy,
1492 Quantity::from("2"),
1493 TimeInForce::Gtc,
1494 UUID4::new(),
1495 0.into(),
1496 false,
1497 false,
1498 None,
1499 None,
1500 None,
1501 None,
1502 Some(ExecAlgorithmId::new("TWAP")),
1503 Some(params),
1504 Some(client_order_id),
1505 None,
1506 ));
1507
1508 let primary_id = order.client_order_id();
1509 let mut submit_params = Params::new();
1510 submit_params.insert(
1511 "routing_profile".to_string(),
1512 serde_json::Value::String("whole-size".to_string()),
1513 );
1514 algo.core
1515 .remember_submit_params(primary_id, Some(submit_params));
1516 algo.on_order(order).unwrap();
1517
1518 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1520 assert_eq!(algo.core.submit_params(&primary_id), None);
1521 }
1522
1523 #[rstest]
1524 fn test_twap_submits_entire_size_when_qty_per_interval_below_min_quantity() {
1525 use nautilus_model::instruments::{InstrumentAny, stubs::crypto_perpetual_ethusdt};
1526
1527 let mut algo = create_twap_algorithm();
1528 register_algorithm(&mut algo);
1529
1530 let mut instrument = crypto_perpetual_ethusdt();
1531 instrument.min_quantity = Some(Quantity::from("0.5"));
1532 {
1533 let cache_rc = algo.core.cache_rc();
1534 let mut cache = cache_rc.borrow_mut();
1535 cache
1536 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
1537 .unwrap();
1538 }
1539
1540 let mut params = IndexMap::new();
1541 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1542 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1543
1544 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1545 let primary_id = order.client_order_id();
1546 let mut submit_params = Params::new();
1547 submit_params.insert(
1548 "routing_profile".to_string(),
1549 serde_json::Value::String("minimum-quantity".to_string()),
1550 );
1551 algo.core
1552 .remember_submit_params(primary_id, Some(submit_params));
1553
1554 algo.on_order(order).unwrap();
1555
1556 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1557 assert_eq!(algo.core.submit_params(&primary_id), None);
1558 }
1559
1560 #[rstest]
1561 fn test_twap_rejects_negative_interval_secs() {
1562 let mut algo = create_twap_algorithm();
1563 register_algorithm(&mut algo);
1564
1565 add_instrument_to_cache(&algo);
1566
1567 let mut params = IndexMap::new();
1568 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1569 params.insert(Ustr::from("interval_secs"), Ustr::from("-0.5"));
1570
1571 let order = create_market_order_with_params(params);
1572 assert_twap_denied(
1573 &mut algo,
1574 &order,
1575 "VALIDATION_FAILED: interval_secs=-0.5 must be finite and positive",
1576 );
1577 }
1578
1579 #[rstest]
1580 fn test_twap_rejects_negative_horizon_secs() {
1581 let mut algo = create_twap_algorithm();
1582 register_algorithm(&mut algo);
1583
1584 add_instrument_to_cache(&algo);
1585
1586 let mut params = IndexMap::new();
1587 params.insert(Ustr::from("horizon_secs"), Ustr::from("-10"));
1588 params.insert(Ustr::from("interval_secs"), Ustr::from("1"));
1589
1590 let order = create_market_order_with_params(params);
1591 assert_twap_denied(
1592 &mut algo,
1593 &order,
1594 "VALIDATION_FAILED: horizon_secs=-10 must be finite and positive",
1595 );
1596 }
1597
1598 #[rstest]
1599 fn test_twap_rejects_zero_interval_secs() {
1600 let mut algo = create_twap_algorithm();
1601 register_algorithm(&mut algo);
1602
1603 add_instrument_to_cache(&algo);
1604
1605 let mut params = IndexMap::new();
1606 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1607 params.insert(Ustr::from("interval_secs"), Ustr::from("0"));
1608
1609 let order = create_market_order_with_params(params);
1610 assert_twap_denied(
1611 &mut algo,
1612 &order,
1613 "VALIDATION_FAILED: interval_secs=0 must be finite and positive",
1614 );
1615 }
1616
1617 #[rstest]
1618 fn test_twap_rejects_huge_finite_interval_before_submission() {
1619 let mut algo = create_twap_algorithm();
1620 register_algorithm(&mut algo);
1621
1622 add_instrument_to_cache(&algo);
1623
1624 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1625 let handler = msgbus::TypedIntoHandler::from({
1626 let captured = received.clone();
1627 move |cmd: TradingCommand| {
1628 if let TradingCommand::SubmitOrder(cmd) = cmd {
1629 *captured.borrow_mut() = Some(cmd);
1630 }
1631 }
1632 });
1633 msgbus::register_trading_command_endpoint(
1634 MessagingSwitchboard::risk_engine_queue_execute(),
1635 handler,
1636 );
1637
1638 let mut params = IndexMap::new();
1639 params.insert(Ustr::from("horizon_secs"), Ustr::from("2e20"));
1640 params.insert(Ustr::from("interval_secs"), Ustr::from("1e20"));
1641
1642 let duration_error = Duration::try_from_secs_f64(1e20).unwrap_err();
1643 let reason = OrderDeniedReason::ValidationFailed {
1644 detail: format!(
1645 "interval_secs=100000000000000000000 is not a valid duration: {duration_error}"
1646 ),
1647 }
1648 .to_string();
1649 assert_twap_denied(&mut algo, &create_market_order_with_params(params), &reason);
1650
1651 assert!(received.borrow().is_none());
1652 }
1653
1654 #[rstest]
1655 fn test_twap_rejects_subnanosecond_interval_before_submission() {
1656 let mut algo = create_twap_algorithm();
1657 register_algorithm(&mut algo);
1658
1659 add_instrument_to_cache(&algo);
1660
1661 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1662 let handler = msgbus::TypedIntoHandler::from({
1663 let captured = received.clone();
1664 move |cmd: TradingCommand| {
1665 if let TradingCommand::SubmitOrder(cmd) = cmd {
1666 *captured.borrow_mut() = Some(cmd);
1667 }
1668 }
1669 });
1670 msgbus::register_trading_command_endpoint(
1671 MessagingSwitchboard::risk_engine_queue_execute(),
1672 handler,
1673 );
1674
1675 let mut params = IndexMap::new();
1676 params.insert(Ustr::from("horizon_secs"), Ustr::from("2e-10"));
1677 params.insert(Ustr::from("interval_secs"), Ustr::from("1e-10"));
1678
1679 assert_twap_denied(
1680 &mut algo,
1681 &create_market_order_with_params(params),
1682 "VALIDATION_FAILED: interval_secs=0.0000000001 rounds to a zero duration",
1683 );
1684
1685 assert!(received.borrow().is_none());
1686 }
1687
1688 #[rstest]
1689 fn test_twap_rejects_interval_exceeding_timestamp_headroom_before_submission() {
1690 let mut algo = create_twap_algorithm();
1691 register_algorithm(&mut algo);
1692
1693 add_instrument_to_cache(&algo);
1694 DataActorNative::clock_mut(&mut algo)
1695 .as_any_mut()
1696 .downcast_mut::<TestClock>()
1697 .unwrap()
1698 .set_time(UnixNanos::new(u64::MAX - 500_000_000));
1699
1700 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1701 let handler = msgbus::TypedIntoHandler::from({
1702 let captured = received.clone();
1703 move |cmd: TradingCommand| {
1704 if let TradingCommand::SubmitOrder(cmd) = cmd {
1705 *captured.borrow_mut() = Some(cmd);
1706 }
1707 }
1708 });
1709 msgbus::register_trading_command_endpoint(
1710 MessagingSwitchboard::risk_engine_queue_execute(),
1711 handler,
1712 );
1713
1714 let mut params = IndexMap::new();
1715 params.insert(Ustr::from("horizon_secs"), Ustr::from("2"));
1716 params.insert(Ustr::from("interval_secs"), Ustr::from("1"));
1717
1718 assert_twap_denied(
1719 &mut algo,
1720 &create_market_order_with_params(params),
1721 "VALIDATION_FAILED: interval_secs=1 exceeds the clock timestamp headroom",
1722 );
1723
1724 assert!(received.borrow().is_none());
1725 }
1726
1727 #[rstest]
1728 fn test_twap_rejects_nan_interval_secs() {
1729 let mut algo = create_twap_algorithm();
1730 register_algorithm(&mut algo);
1731
1732 add_instrument_to_cache(&algo);
1733
1734 let mut params = IndexMap::new();
1735 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1736 params.insert(Ustr::from("interval_secs"), Ustr::from("NaN"));
1737
1738 let order = create_market_order_with_params(params);
1739 assert_twap_denied(
1740 &mut algo,
1741 &order,
1742 "VALIDATION_FAILED: interval_secs=NaN must be finite and positive",
1743 );
1744 }
1745
1746 #[rstest]
1747 fn test_twap_rejects_infinity_horizon_secs() {
1748 let mut algo = create_twap_algorithm();
1749 register_algorithm(&mut algo);
1750
1751 add_instrument_to_cache(&algo);
1752
1753 let mut params = IndexMap::new();
1754 params.insert(Ustr::from("horizon_secs"), Ustr::from("inf"));
1755 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
1756
1757 let order = create_market_order_with_params(params);
1758 assert_twap_denied(
1759 &mut algo,
1760 &order,
1761 "VALIDATION_FAILED: horizon_secs=inf must be finite and positive",
1762 );
1763 }
1764}