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::{ModifyOrder, 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_refused_modify_preserves_remaining_quantity_schedule() {
934 let mut algo = create_twap_algorithm();
935 register_algorithm(&mut algo);
936 add_instrument_to_cache(&algo);
937
938 let mut params = IndexMap::new();
939 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
940 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
941 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
942 let primary_id = order.client_order_id();
943 algo.on_order(order).unwrap();
944
945 let command = ModifyOrder::new(
946 TraderId::from("TRADER-001"),
947 None,
948 StrategyId::from("STRAT-001"),
949 InstrumentId::from("ETHUSDT-PERP.BINANCE"),
950 primary_id,
951 None,
952 Some(Quantity::from("0.4")),
953 None,
954 None,
955 UUID4::new(),
956 0.into(),
957 None,
958 None,
959 );
960
961 algo.handle_modify_order(command).unwrap();
962
963 let primary_quantity = algo.cache().order(&primary_id).unwrap().quantity();
964 let scheduled_quantity = algo.scheduled_orders[&primary_id]
965 .remaining_sizes
966 .iter()
967 .fold(Decimal::ZERO, |total, quantity| {
968 total + quantity.as_decimal()
969 });
970 assert_eq!(primary_quantity.as_decimal(), scheduled_quantity);
971 assert_eq!(primary_quantity, Quantity::from("0.8"));
972 }
973
974 #[rstest]
975 fn test_twap_on_order_accepts_already_cached_primary() {
976 let mut algo = create_twap_algorithm();
977 register_algorithm(&mut algo);
978
979 add_instrument_to_cache(&algo);
980
981 let mut params = IndexMap::new();
982 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
983 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
984
985 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
986 let primary_id = order.client_order_id();
987
988 {
990 let cache_rc = algo.core.cache_rc();
991 let mut cache = cache_rc.borrow_mut();
992 cache.add_order(order.clone(), None, None, false).unwrap();
993 }
994
995 algo.on_order(order).unwrap();
996
997 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 2);
998 }
999
1000 #[rstest]
1001 fn test_twap_calculates_size_schedule_with_remainder() {
1002 let mut algo = create_twap_algorithm();
1003 register_algorithm(&mut algo);
1004
1005 add_instrument_to_cache(&algo);
1006 let instrument = nautilus_model::instruments::stubs::crypto_perpetual_ethusdt();
1007
1008 let mut params = IndexMap::new();
1010 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1011 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1012
1013 let order = create_market_order_with_params(params);
1014 let primary_id = order.client_order_id();
1015
1016 algo.on_order(order).unwrap();
1017
1018 let remaining = &algo
1020 .scheduled_orders
1021 .get(&primary_id)
1022 .unwrap()
1023 .remaining_sizes;
1024 assert_eq!(remaining.len(), 3);
1025 assert_eq!(
1026 remaining,
1027 &[
1028 Quantity::from("0.333"),
1029 Quantity::from("0.333"),
1030 Quantity::from("0.001"),
1031 ]
1032 );
1033
1034 for quantity in remaining {
1035 assert_eq!(quantity.precision, instrument.size_precision());
1036 assert_eq!(quantity.raw % instrument.size_increment().raw, 0);
1037 }
1038
1039 let first = algo
1040 .cache()
1041 .order(&ClientOrderId::from("O-001-E1"))
1042 .unwrap()
1043 .quantity();
1044 let total = remaining
1045 .iter()
1046 .fold(first.as_decimal(), |sum, qty| sum + qty.as_decimal());
1047 assert_eq!(total, Quantity::from("1.0").as_decimal());
1048 }
1049
1050 #[rstest]
1051 fn test_twap_children_use_instrument_size_precision() {
1052 let mut algo = create_twap_algorithm();
1053 register_algorithm(&mut algo);
1054
1055 add_instrument_to_cache(&algo);
1056 let instrument = nautilus_model::instruments::stubs::crypto_perpetual_ethusdt();
1057
1058 let mut params = IndexMap::new();
1059 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1060 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1061
1062 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.000000000"));
1063 let primary_id = order.client_order_id();
1064
1065 algo.on_order(order).unwrap();
1066
1067 let spawned = algo
1068 .cache()
1069 .order(&ClientOrderId::from("O-001-E1"))
1070 .unwrap()
1071 .quantity();
1072 assert_eq!(spawned.precision, instrument.size_precision());
1073 assert!(
1074 algo.scheduled_orders[&primary_id]
1075 .remaining_sizes
1076 .iter()
1077 .all(|quantity| quantity.precision == instrument.size_precision())
1078 );
1079 }
1080
1081 #[rstest]
1082 fn test_twap_on_time_event_spawns_next_slice() {
1083 let mut algo = create_twap_algorithm();
1084 register_algorithm(&mut algo);
1085
1086 add_instrument_to_cache(&algo);
1087
1088 let mut params = IndexMap::new();
1090 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1091 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1092
1093 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1094 let primary_id = order.client_order_id();
1095 let mut submit_params = Params::new();
1096 submit_params.insert(
1097 "routing_profile".to_string(),
1098 serde_json::Value::String("intermediate-slice".to_string()),
1099 );
1100 algo.core
1101 .remember_submit_params(primary_id, Some(submit_params.clone()));
1102
1103 algo.on_order(order).unwrap();
1104
1105 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 2);
1107
1108 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1110 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1111
1112 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 1);
1114 assert_eq!(algo.core.submit_params(&primary_id), Some(submit_params));
1115 }
1116
1117 #[rstest]
1118 fn test_twap_data_actor_dispatch_spawns_next_slice() {
1119 let mut algo = create_twap_algorithm();
1120 register_algorithm(&mut algo);
1121
1122 add_instrument_to_cache(&algo);
1123
1124 let mut params = IndexMap::new();
1125 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1126 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1127
1128 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1129 let primary_id = order.client_order_id();
1130
1131 algo.on_order(order).unwrap();
1132 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 2);
1133
1134 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1136 algo.handle_time_event(&event);
1137
1138 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 1);
1139 }
1140
1141 #[rstest]
1142 fn test_twap_on_time_event_completes_on_final_slice() {
1143 let mut algo = create_twap_algorithm();
1144 register_algorithm(&mut algo);
1145
1146 add_instrument_to_cache(&algo);
1147
1148 let mut params = IndexMap::new();
1150 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1151 params.insert(Ustr::from("interval_secs"), Ustr::from("30"));
1152
1153 let order = create_market_order_with_params(params);
1154 let primary_id = order.client_order_id();
1155 let mut submit_params = Params::new();
1156 submit_params.insert(
1157 "routing_profile".to_string(),
1158 serde_json::Value::String("final-slice".to_string()),
1159 );
1160 algo.core
1161 .remember_submit_params(primary_id, Some(submit_params.clone()));
1162
1163 algo.on_order(order).unwrap();
1164 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 1);
1165 assert_eq!(
1166 algo.core.submit_params(&primary_id),
1167 Some(submit_params.clone())
1168 );
1169
1170 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1171 let handler = msgbus::TypedIntoHandler::from({
1172 let captured = received.clone();
1173 move |cmd: TradingCommand| {
1174 if let TradingCommand::SubmitOrder(cmd) = cmd {
1175 *captured.borrow_mut() = Some(cmd);
1176 }
1177 }
1178 });
1179 msgbus::register_trading_command_endpoint(
1180 MessagingSwitchboard::risk_engine_queue_execute(),
1181 handler,
1182 );
1183
1184 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1186 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1187
1188 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1190 assert_eq!(
1191 received
1192 .borrow()
1193 .as_ref()
1194 .and_then(|cmd| cmd.params.clone()),
1195 Some(submit_params),
1196 );
1197 assert_eq!(algo.core.submit_params(&primary_id), None);
1198 }
1199
1200 #[rstest]
1201 fn test_twap_on_time_event_completes_when_primary_closed() {
1202 let mut algo = create_twap_algorithm();
1203 register_algorithm(&mut algo);
1204
1205 add_instrument_to_cache(&algo);
1206
1207 let mut params = IndexMap::new();
1208 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1209 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1210
1211 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1212 let primary_id = order.client_order_id();
1213 let mut submit_params = Params::new();
1214 submit_params.insert(
1215 "routing_profile".to_string(),
1216 serde_json::Value::String("closed-primary".to_string()),
1217 );
1218 algo.core
1219 .remember_submit_params(primary_id, Some(submit_params));
1220
1221 algo.on_order(order).unwrap();
1222 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 2);
1223
1224 {
1226 let cache_rc = algo.core.cache_rc();
1227 let mut cache = cache_rc.borrow_mut();
1228 let primary = cache.order(&primary_id).map(|o| o.clone()).unwrap();
1229
1230 let canceled = OrderCanceledSpec::builder()
1231 .trader_id(primary.trader_id())
1232 .strategy_id(primary.strategy_id())
1233 .instrument_id(primary.instrument_id())
1234 .client_order_id(primary.client_order_id())
1235 .build();
1236 cache
1237 .update_order(&OrderEventAny::Canceled(canceled))
1238 .unwrap();
1239 }
1240
1241 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1243 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1244
1245 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1247 assert_eq!(algo.core.submit_params(&primary_id), None);
1248 }
1249
1250 #[rstest]
1251 fn test_twap_on_time_event_completes_when_primary_missing() {
1252 let mut algo = create_twap_algorithm();
1253 register_algorithm(&mut algo);
1254 add_instrument_to_cache(&algo);
1255
1256 let mut params = IndexMap::new();
1257 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1258 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1259
1260 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1261 let primary_id = order.client_order_id();
1262 algo.on_order(order).unwrap();
1263
1264 {
1265 let cache_rc = algo.core.cache_rc();
1266 let mut cache = cache_rc.borrow_mut();
1267 let primary = cache.order(&primary_id).map(|order| order.clone()).unwrap();
1268 let canceled = OrderCanceledSpec::builder()
1269 .trader_id(primary.trader_id())
1270 .strategy_id(primary.strategy_id())
1271 .instrument_id(primary.instrument_id())
1272 .client_order_id(primary.client_order_id())
1273 .build();
1274 cache
1275 .update_order(&OrderEventAny::Canceled(canceled))
1276 .unwrap();
1277 cache.purge_order(primary_id);
1278 }
1279 assert!(algo.cache().order(&primary_id).is_none());
1280
1281 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1282 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1283
1284 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1286 assert!(algo.clock().timer_names().is_empty());
1287 }
1288
1289 #[rstest]
1290 fn test_twap_on_stop_cancels_timers() {
1291 let mut algo = create_twap_algorithm();
1292 register_algorithm(&mut algo);
1293
1294 add_instrument_to_cache(&algo);
1295
1296 let mut params = IndexMap::new();
1297 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1298 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1299
1300 let order = create_market_order_with_params(params);
1301 let primary_id = order.client_order_id();
1302
1303 algo.on_order(order).unwrap();
1304
1305 assert!(
1307 algo.clock()
1308 .timer_names()
1309 .iter()
1310 .any(|name| name.as_str() == primary_id.as_str())
1311 );
1312
1313 DataActor::on_stop(&mut algo).unwrap();
1315
1316 assert!(algo.clock().timer_names().is_empty());
1318 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 3);
1319 }
1320
1321 #[rstest]
1322 fn test_twap_on_resume_rearms_timer_without_submitting() {
1323 let mut algo = create_twap_algorithm();
1324 register_algorithm(&mut algo);
1325 add_instrument_to_cache(&algo);
1326
1327 let mut params = IndexMap::new();
1328 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1329 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1330
1331 let order = create_market_order_with_params(params);
1332 let primary_id = order.client_order_id();
1333 algo.on_order(order).unwrap();
1334 let order_count = algo
1335 .cache()
1336 .orders_total_count(None, None, None, None, None);
1337
1338 Component::stop(&mut algo).unwrap();
1339
1340 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1343 let handler = msgbus::TypedIntoHandler::from({
1344 let captured = received.clone();
1345 move |cmd: TradingCommand| {
1346 if let TradingCommand::SubmitOrder(cmd) = cmd {
1347 *captured.borrow_mut() = Some(cmd);
1348 }
1349 }
1350 });
1351 msgbus::register_trading_command_endpoint(
1352 MessagingSwitchboard::risk_engine_queue_execute(),
1353 handler,
1354 );
1355
1356 let resume_time = algo.clock().timestamp_ns();
1357 Component::resume(&mut algo).unwrap();
1358
1359 assert!(received.borrow().is_none());
1360 assert_eq!(algo.clock().timer_count(), 1);
1361 assert_eq!(
1362 algo.clock().next_time_ns(primary_id.as_str()),
1363 Some(resume_time + 20_000_000_000)
1364 );
1365 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 3);
1366 assert_eq!(
1367 algo.cache()
1368 .orders_total_count(None, None, None, None, None),
1369 order_count
1370 );
1371 }
1372
1373 #[rstest]
1374 fn test_twap_on_resume_executes_remaining_slices() {
1375 let mut algo = create_twap_algorithm();
1376 let clock = register_algorithm_with_clock(&mut algo);
1377 add_instrument_to_cache(&algo);
1378
1379 let mut params = IndexMap::new();
1380 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1381 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1382
1383 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1384 let primary_id = order.client_order_id();
1385 algo.on_order(order).unwrap();
1386
1387 Component::stop(&mut algo).unwrap();
1388 Component::resume(&mut algo).unwrap();
1389 assert_eq!(algo.clock().timer_count(), 1);
1390
1391 let first_events = clock.borrow_mut().advance_time(20_000_000_000.into(), true);
1392 assert_eq!(first_events.len(), 1);
1393 algo.handle_time_event(&first_events[0]);
1394 assert_eq!(algo.scheduled_orders[&primary_id].remaining_sizes.len(), 1);
1395
1396 let final_events = clock.borrow_mut().advance_time(40_000_000_000.into(), true);
1397 assert_eq!(final_events.len(), 1);
1398 algo.handle_time_event(&final_events[0]);
1399 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1400 }
1401
1402 #[rstest]
1403 fn test_twap_on_resume_completes_closed_primary() {
1404 let mut algo = create_twap_algorithm();
1405 register_algorithm(&mut algo);
1406 add_instrument_to_cache(&algo);
1407
1408 let mut params = IndexMap::new();
1409 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1410 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1411
1412 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1413 let primary_id = order.client_order_id();
1414 algo.on_order(order).unwrap();
1415 Component::stop(&mut algo).unwrap();
1416
1417 {
1418 let cache_rc = algo.core.cache_rc();
1419 let mut cache = cache_rc.borrow_mut();
1420 let primary = cache.order(&primary_id).map(|order| order.clone()).unwrap();
1421 let canceled = OrderCanceledSpec::builder()
1422 .trader_id(primary.trader_id())
1423 .strategy_id(primary.strategy_id())
1424 .instrument_id(primary.instrument_id())
1425 .client_order_id(primary.client_order_id())
1426 .build();
1427 cache
1428 .update_order(&OrderEventAny::Canceled(canceled))
1429 .unwrap();
1430 }
1431
1432 Component::resume(&mut algo).unwrap();
1433
1434 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1435 assert_eq!(algo.clock().timer_count(), 0);
1436 }
1437
1438 #[rstest]
1439 fn test_twap_on_resume_completes_missing_primary() {
1440 let mut algo = create_twap_algorithm();
1441 register_algorithm(&mut algo);
1442 add_instrument_to_cache(&algo);
1443
1444 let mut params = IndexMap::new();
1445 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1446 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1447
1448 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1449 let primary_id = order.client_order_id();
1450 algo.on_order(order).unwrap();
1451 Component::stop(&mut algo).unwrap();
1452
1453 {
1454 let cache_rc = algo.core.cache_rc();
1455 let mut cache = cache_rc.borrow_mut();
1456 let primary = cache.order(&primary_id).map(|order| order.clone()).unwrap();
1457 let canceled = OrderCanceledSpec::builder()
1458 .trader_id(primary.trader_id())
1459 .strategy_id(primary.strategy_id())
1460 .instrument_id(primary.instrument_id())
1461 .client_order_id(primary.client_order_id())
1462 .build();
1463 cache
1464 .update_order(&OrderEventAny::Canceled(canceled))
1465 .unwrap();
1466 cache.purge_order(primary_id);
1467 }
1468
1469 assert!(algo.cache().order(&primary_id).is_none());
1470 Component::resume(&mut algo).unwrap();
1471
1472 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1473 assert_eq!(algo.clock().timer_count(), 0);
1474 }
1475
1476 #[rstest]
1477 fn test_twap_fractional_interval_secs() {
1478 let mut algo = create_twap_algorithm();
1479 register_algorithm(&mut algo);
1480
1481 add_instrument_to_cache(&algo);
1482
1483 let mut params = IndexMap::new();
1485 params.insert(Ustr::from("horizon_secs"), Ustr::from("3"));
1486 params.insert(Ustr::from("interval_secs"), Ustr::from("0.5"));
1487
1488 let order = create_market_order_with_params(params);
1489 let primary_id = order.client_order_id();
1490
1491 algo.on_order(order).unwrap();
1493
1494 let remaining = &algo
1496 .scheduled_orders
1497 .get(&primary_id)
1498 .unwrap()
1499 .remaining_sizes;
1500 assert!(remaining.len() >= 5);
1501 }
1502
1503 #[rstest]
1504 fn test_twap_submits_entire_size_when_qty_per_interval_below_size_increment() {
1505 use nautilus_model::instruments::{InstrumentAny, stubs::equity_aapl};
1506
1507 let mut algo = create_twap_algorithm();
1508 register_algorithm(&mut algo);
1509
1510 let instrument = equity_aapl();
1512 let instrument_id = instrument.id();
1513 {
1514 let cache_rc = algo.core.cache_rc();
1515 let mut cache = cache_rc.borrow_mut();
1516 cache
1517 .add_instrument(InstrumentAny::Equity(instrument))
1518 .unwrap();
1519 }
1520
1521 let mut params = IndexMap::new();
1524 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1525 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
1526
1527 let client_order_id = ClientOrderId::from("O-002");
1528 let order = OrderAny::Market(MarketOrder::new(
1529 TraderId::from("TRADER-001"),
1530 StrategyId::from("STRAT-001"),
1531 instrument_id,
1532 client_order_id,
1533 OrderSide::Buy,
1534 Quantity::from("2"),
1535 TimeInForce::Gtc,
1536 UUID4::new(),
1537 0.into(),
1538 false,
1539 false,
1540 None,
1541 None,
1542 None,
1543 None,
1544 Some(ExecAlgorithmId::new("TWAP")),
1545 Some(params),
1546 Some(client_order_id),
1547 None,
1548 ));
1549
1550 let primary_id = order.client_order_id();
1551 let mut submit_params = Params::new();
1552 submit_params.insert(
1553 "routing_profile".to_string(),
1554 serde_json::Value::String("whole-size".to_string()),
1555 );
1556 algo.core
1557 .remember_submit_params(primary_id, Some(submit_params));
1558 algo.on_order(order).unwrap();
1559
1560 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1562 assert_eq!(algo.core.submit_params(&primary_id), None);
1563 }
1564
1565 #[rstest]
1566 fn test_twap_submits_entire_size_when_qty_per_interval_below_min_quantity() {
1567 use nautilus_model::instruments::{InstrumentAny, stubs::crypto_perpetual_ethusdt};
1568
1569 let mut algo = create_twap_algorithm();
1570 register_algorithm(&mut algo);
1571
1572 let mut instrument = crypto_perpetual_ethusdt();
1573 instrument.min_quantity = Some(Quantity::from("0.5"));
1574 {
1575 let cache_rc = algo.core.cache_rc();
1576 let mut cache = cache_rc.borrow_mut();
1577 cache
1578 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
1579 .unwrap();
1580 }
1581
1582 let mut params = IndexMap::new();
1583 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1584 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1585
1586 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1587 let primary_id = order.client_order_id();
1588 let mut submit_params = Params::new();
1589 submit_params.insert(
1590 "routing_profile".to_string(),
1591 serde_json::Value::String("minimum-quantity".to_string()),
1592 );
1593 algo.core
1594 .remember_submit_params(primary_id, Some(submit_params));
1595
1596 algo.on_order(order).unwrap();
1597
1598 assert!(algo.scheduled_orders.get(&primary_id).is_none());
1599 assert_eq!(algo.core.submit_params(&primary_id), None);
1600 }
1601
1602 #[rstest]
1603 fn test_twap_rejects_negative_interval_secs() {
1604 let mut algo = create_twap_algorithm();
1605 register_algorithm(&mut algo);
1606
1607 add_instrument_to_cache(&algo);
1608
1609 let mut params = IndexMap::new();
1610 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1611 params.insert(Ustr::from("interval_secs"), Ustr::from("-0.5"));
1612
1613 let order = create_market_order_with_params(params);
1614 assert_twap_denied(
1615 &mut algo,
1616 &order,
1617 "VALIDATION_FAILED: interval_secs=-0.5 must be finite and positive",
1618 );
1619 }
1620
1621 #[rstest]
1622 fn test_twap_rejects_negative_horizon_secs() {
1623 let mut algo = create_twap_algorithm();
1624 register_algorithm(&mut algo);
1625
1626 add_instrument_to_cache(&algo);
1627
1628 let mut params = IndexMap::new();
1629 params.insert(Ustr::from("horizon_secs"), Ustr::from("-10"));
1630 params.insert(Ustr::from("interval_secs"), Ustr::from("1"));
1631
1632 let order = create_market_order_with_params(params);
1633 assert_twap_denied(
1634 &mut algo,
1635 &order,
1636 "VALIDATION_FAILED: horizon_secs=-10 must be finite and positive",
1637 );
1638 }
1639
1640 #[rstest]
1641 fn test_twap_rejects_zero_interval_secs() {
1642 let mut algo = create_twap_algorithm();
1643 register_algorithm(&mut algo);
1644
1645 add_instrument_to_cache(&algo);
1646
1647 let mut params = IndexMap::new();
1648 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1649 params.insert(Ustr::from("interval_secs"), Ustr::from("0"));
1650
1651 let order = create_market_order_with_params(params);
1652 assert_twap_denied(
1653 &mut algo,
1654 &order,
1655 "VALIDATION_FAILED: interval_secs=0 must be finite and positive",
1656 );
1657 }
1658
1659 #[rstest]
1660 fn test_twap_rejects_huge_finite_interval_before_submission() {
1661 let mut algo = create_twap_algorithm();
1662 register_algorithm(&mut algo);
1663
1664 add_instrument_to_cache(&algo);
1665
1666 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1667 let handler = msgbus::TypedIntoHandler::from({
1668 let captured = received.clone();
1669 move |cmd: TradingCommand| {
1670 if let TradingCommand::SubmitOrder(cmd) = cmd {
1671 *captured.borrow_mut() = Some(cmd);
1672 }
1673 }
1674 });
1675 msgbus::register_trading_command_endpoint(
1676 MessagingSwitchboard::risk_engine_queue_execute(),
1677 handler,
1678 );
1679
1680 let mut params = IndexMap::new();
1681 params.insert(Ustr::from("horizon_secs"), Ustr::from("2e20"));
1682 params.insert(Ustr::from("interval_secs"), Ustr::from("1e20"));
1683
1684 let duration_error = Duration::try_from_secs_f64(1e20).unwrap_err();
1685 let reason = OrderDeniedReason::ValidationFailed {
1686 detail: format!(
1687 "interval_secs=100000000000000000000 is not a valid duration: {duration_error}"
1688 ),
1689 }
1690 .to_string();
1691 assert_twap_denied(&mut algo, &create_market_order_with_params(params), &reason);
1692
1693 assert!(received.borrow().is_none());
1694 }
1695
1696 #[rstest]
1697 fn test_twap_rejects_subnanosecond_interval_before_submission() {
1698 let mut algo = create_twap_algorithm();
1699 register_algorithm(&mut algo);
1700
1701 add_instrument_to_cache(&algo);
1702
1703 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1704 let handler = msgbus::TypedIntoHandler::from({
1705 let captured = received.clone();
1706 move |cmd: TradingCommand| {
1707 if let TradingCommand::SubmitOrder(cmd) = cmd {
1708 *captured.borrow_mut() = Some(cmd);
1709 }
1710 }
1711 });
1712 msgbus::register_trading_command_endpoint(
1713 MessagingSwitchboard::risk_engine_queue_execute(),
1714 handler,
1715 );
1716
1717 let mut params = IndexMap::new();
1718 params.insert(Ustr::from("horizon_secs"), Ustr::from("2e-10"));
1719 params.insert(Ustr::from("interval_secs"), Ustr::from("1e-10"));
1720
1721 assert_twap_denied(
1722 &mut algo,
1723 &create_market_order_with_params(params),
1724 "VALIDATION_FAILED: interval_secs=0.0000000001 rounds to a zero duration",
1725 );
1726
1727 assert!(received.borrow().is_none());
1728 }
1729
1730 #[rstest]
1731 fn test_twap_rejects_interval_exceeding_timestamp_headroom_before_submission() {
1732 let mut algo = create_twap_algorithm();
1733 register_algorithm(&mut algo);
1734
1735 add_instrument_to_cache(&algo);
1736 DataActorNative::clock_mut(&mut algo)
1737 .as_any_mut()
1738 .downcast_mut::<TestClock>()
1739 .unwrap()
1740 .set_time(UnixNanos::new(u64::MAX - 500_000_000));
1741
1742 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1743 let handler = msgbus::TypedIntoHandler::from({
1744 let captured = received.clone();
1745 move |cmd: TradingCommand| {
1746 if let TradingCommand::SubmitOrder(cmd) = cmd {
1747 *captured.borrow_mut() = Some(cmd);
1748 }
1749 }
1750 });
1751 msgbus::register_trading_command_endpoint(
1752 MessagingSwitchboard::risk_engine_queue_execute(),
1753 handler,
1754 );
1755
1756 let mut params = IndexMap::new();
1757 params.insert(Ustr::from("horizon_secs"), Ustr::from("2"));
1758 params.insert(Ustr::from("interval_secs"), Ustr::from("1"));
1759
1760 assert_twap_denied(
1761 &mut algo,
1762 &create_market_order_with_params(params),
1763 "VALIDATION_FAILED: interval_secs=1 exceeds the clock timestamp headroom",
1764 );
1765
1766 assert!(received.borrow().is_none());
1767 }
1768
1769 #[rstest]
1770 fn test_twap_rejects_nan_interval_secs() {
1771 let mut algo = create_twap_algorithm();
1772 register_algorithm(&mut algo);
1773
1774 add_instrument_to_cache(&algo);
1775
1776 let mut params = IndexMap::new();
1777 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1778 params.insert(Ustr::from("interval_secs"), Ustr::from("NaN"));
1779
1780 let order = create_market_order_with_params(params);
1781 assert_twap_denied(
1782 &mut algo,
1783 &order,
1784 "VALIDATION_FAILED: interval_secs=NaN must be finite and positive",
1785 );
1786 }
1787
1788 #[rstest]
1789 fn test_twap_rejects_infinity_horizon_secs() {
1790 let mut algo = create_twap_algorithm();
1791 register_algorithm(&mut algo);
1792
1793 add_instrument_to_cache(&algo);
1794
1795 let mut params = IndexMap::new();
1796 params.insert(Ustr::from("horizon_secs"), Ustr::from("inf"));
1797 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
1798
1799 let order = create_market_order_with_params(params);
1800 assert_twap_denied(
1801 &mut algo,
1802 &order,
1803 "VALIDATION_FAILED: horizon_secs=inf must be finite and positive",
1804 );
1805 }
1806}