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_sizes: AHashMap<ClientOrderId, Vec<Quantity>>,
70}
71
72impl TwapAlgorithm {
73 #[must_use]
75 pub fn new(config: TwapAlgorithmConfig) -> Self {
76 Self {
77 core: ExecutionAlgorithmCore::new(config),
78 scheduled_sizes: 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_sizes.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_reset(&mut self) -> anyhow::Result<()> {
107 ExecutionAlgorithm::on_reset(self)
108 }
109}
110
111nautilus_execution_algorithm!(TwapAlgorithm, {
112 fn on_order(&mut self, order: OrderAny) -> anyhow::Result<()> {
113 let primary_id = order.client_order_id();
114
115 if self.scheduled_sizes.contains_key(&primary_id) {
116 anyhow::bail!("Order {primary_id} already being executed");
117 }
118
119 log::info!("Received order for TWAP execution: {order:?}");
120
121 if order.order_type() != OrderType::Market {
123 let reason = OrderDeniedReason::UnsupportedOrderType {
124 order_type: order.order_type(),
125 }
126 .to_string();
127 return self.deny_order(&order, Ustr::from(&reason));
128 }
129
130 let instrument = {
131 let cache = ExecutionAlgorithmNative::exec_algorithm_core(self).cache_ref();
132 cache.instrument(&order.instrument_id()).cloned()
133 };
134
135 let Some(instrument) = instrument else {
136 let reason = OrderDeniedReason::InstrumentNotFound {
137 instrument_id: order.instrument_id(),
138 }
139 .to_string();
140 return self.deny_order(&order, Ustr::from(&reason));
141 };
142
143 let Some(exec_params) = order.exec_algorithm_params() else {
144 return self.deny_order(&order, validation_failed("exec_algorithm_params not found"));
145 };
146
147 let Some(horizon_secs_str) = exec_params.get(&Ustr::from("horizon_secs")) else {
148 return self.deny_order(
149 &order,
150 validation_failed("horizon_secs not found in exec_algorithm_params"),
151 );
152 };
153
154 let horizon_secs: f64 = match horizon_secs_str.parse() {
155 Ok(value) => value,
156 Err(_) => {
157 return self.deny_order(
158 &order,
159 validation_failed(format!(
160 "horizon_secs={horizon_secs_str} is not a valid number"
161 )),
162 );
163 }
164 };
165
166 let Some(interval_secs_str) = exec_params.get(&Ustr::from("interval_secs")) else {
167 return self.deny_order(
168 &order,
169 validation_failed("interval_secs not found in exec_algorithm_params"),
170 );
171 };
172
173 let interval_secs: f64 = match interval_secs_str.parse() {
174 Ok(value) => value,
175 Err(_) => {
176 return self.deny_order(
177 &order,
178 validation_failed(format!(
179 "interval_secs={interval_secs_str} is not a valid number"
180 )),
181 );
182 }
183 };
184
185 if !horizon_secs.is_finite() || horizon_secs <= 0.0 {
186 return self.deny_order(
187 &order,
188 validation_failed(format!(
189 "horizon_secs={horizon_secs} must be finite and positive"
190 )),
191 );
192 }
193
194 if !interval_secs.is_finite() || interval_secs <= 0.0 {
195 return self.deny_order(
196 &order,
197 validation_failed(format!(
198 "interval_secs={interval_secs} must be finite and positive"
199 )),
200 );
201 }
202
203 if horizon_secs < interval_secs {
204 return self.deny_order(
205 &order,
206 validation_failed(format!(
207 "horizon_secs={horizon_secs} must be greater than or equal to interval_secs={interval_secs}"
208 )),
209 );
210 }
211
212 let num_intervals = (horizon_secs / interval_secs).floor() as u64;
213 if num_intervals == 0 {
214 return self.deny_order(&order, validation_failed("num_intervals is 0"));
215 }
216
217 let total_qty = order.quantity();
218 let interval_count = Decimal::from(num_intervals);
219 let quotient = total_qty.as_decimal() / interval_count;
220 let floored = quotient.round_dp_with_strategy(
221 u32::from(instrument.size_precision()),
222 RoundingStrategy::ToZero,
223 );
224 let qty_per_interval = match instrument.try_make_qty_from_decimal(floored, None) {
225 Ok(quantity) => quantity,
226 Err(e) => {
227 return self.deny_order(
228 &order,
229 validation_failed(format!("invalid qty_per_interval={floored}: {e}")),
230 );
231 }
232 };
233 let remainder = total_qty.as_decimal() - floored * interval_count;
234
235 if qty_per_interval == total_qty || qty_per_interval < instrument.size_increment() {
236 log::warn!(
237 "Submitting for entire size: qty_per_interval={qty_per_interval}, order_quantity={total_qty}"
238 );
239 self.submit_order(order, None, None)?;
240 self.complete_sequence(primary_id);
241 return Ok(());
242 }
243
244 if let Some(min_qty) = instrument.min_quantity()
245 && qty_per_interval < min_qty
246 {
247 log::warn!(
248 "Submitting for entire size: qty_per_interval={qty_per_interval} < min_quantity={min_qty}"
249 );
250 self.submit_order(order, None, None)?;
251 self.complete_sequence(primary_id);
252 return Ok(());
253 }
254
255 let interval = match Duration::try_from_secs_f64(interval_secs) {
256 Ok(interval) => interval,
257 Err(e) => {
258 return self.deny_order(
259 &order,
260 validation_failed(format!(
261 "interval_secs={interval_secs} is not a valid duration: {e}"
262 )),
263 );
264 }
265 };
266
267 if interval == Duration::ZERO {
268 return self.deny_order(
269 &order,
270 validation_failed(format!(
271 "interval_secs={interval_secs} rounds to a zero duration"
272 )),
273 );
274 }
275 let Ok(interval_ns) = u64::try_from(interval.as_nanos()) else {
276 return self.deny_order(
277 &order,
278 validation_failed(format!(
279 "interval_secs={interval_secs} exceeds the clock nanosecond range"
280 )),
281 );
282 };
283 let timestamp_ns = ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
284 .clock_mut()
285 .timestamp_ns()
286 .as_u64();
287
288 if timestamp_ns.checked_add(interval_ns).is_none() {
289 return self.deny_order(
290 &order,
291 validation_failed(format!(
292 "interval_secs={interval_secs} exceeds the clock timestamp headroom"
293 )),
294 );
295 }
296
297 let mut scheduled_sizes: Vec<Quantity> = vec![qty_per_interval; num_intervals as usize];
298
299 if remainder > Decimal::ZERO {
300 let remainder_qty = match instrument.try_make_qty_from_decimal(remainder, None) {
301 Ok(quantity) => quantity,
302 Err(e) => {
303 return self.deny_order(
304 &order,
305 validation_failed(format!("invalid qty_remainder={remainder}: {e}")),
306 );
307 }
308 };
309 scheduled_sizes.push(remainder_qty);
310 }
311
312 let scheduled_total = scheduled_sizes
313 .iter()
314 .fold(Decimal::ZERO, |total, quantity| {
315 total + quantity.as_decimal()
316 });
317
318 if scheduled_total != total_qty.as_decimal() {
319 return self.deny_order(
320 &order,
321 validation_failed(format!(
322 "scheduled quantity {scheduled_total} does not equal order quantity {total_qty}"
323 )),
324 );
325 }
326
327 log::info!("Order execution size schedule: {scheduled_sizes:?}");
328
329 {
332 let cache_rc = ExecutionAlgorithmNative::exec_algorithm_core(self).cache_rc();
333 let mut cache = cache_rc.borrow_mut();
334 if !cache.order_exists(&primary_id) {
335 cache.add_order(order.clone(), None, None, false)?;
336 }
337 }
338
339 self.scheduled_sizes
340 .insert(primary_id, scheduled_sizes.clone());
341
342 let first_qty = self.scheduled_sizes.get_mut(&primary_id).unwrap().remove(0);
343 let is_single_slice = self
344 .scheduled_sizes
345 .get(&primary_id)
346 .is_some_and(|s| s.is_empty());
347
348 if is_single_slice {
350 self.submit_order(order, None, None)?;
351 self.complete_sequence(primary_id);
352 return Ok(());
353 }
354
355 let tags = order.tags().map(|t| t.to_vec());
357 let time_in_force = order.time_in_force();
358 let reduce_only = order.is_reduce_only();
359 let mut order = order;
360 let spawned = self.spawn_market(
361 &mut order,
362 first_qty,
363 time_in_force,
364 reduce_only,
365 tags,
366 true,
367 );
368 self.submit_order(spawned.into(), None, None)?;
369
370 ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
371 .clock_mut()
372 .set_timer(primary_id.as_str(), interval, None, None, None, None, None)?;
373
374 log::info!(
375 "Started TWAP execution for {primary_id}: horizon_secs={horizon_secs}, interval_secs={interval_secs}"
376 );
377
378 Ok(())
379 }
380
381 fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()> {
382 log::info!("Received time event: {event:?}");
383
384 let primary_id = ClientOrderId::new(event.name.as_str());
385
386 let primary = {
387 let cache = ExecutionAlgorithmNative::exec_algorithm_core(self).cache_ref();
388 cache.order(&primary_id).map(|o| o.clone())
389 };
390
391 let Some(primary) = primary else {
392 log::error!("Cannot find primary order for exec_spawn_id={primary_id}");
393 return Ok(());
394 };
395
396 if primary.is_closed() {
397 self.complete_sequence(primary_id);
398 return Ok(());
399 }
400
401 let Some(scheduled_sizes) = self.scheduled_sizes.get_mut(&primary_id) else {
402 log::error!("Cannot find scheduled sizes for exec_spawn_id={primary_id}");
403 return Ok(());
404 };
405
406 if scheduled_sizes.is_empty() {
407 log::warn!("No more size to execute for exec_spawn_id={primary_id}");
408 return Ok(());
409 }
410
411 let quantity = scheduled_sizes.remove(0);
412 let is_final_slice = scheduled_sizes.is_empty();
413
414 if is_final_slice {
416 self.submit_order(primary, None, None)?;
417 self.complete_sequence(primary_id);
418 return Ok(());
419 }
420
421 let tags = primary.tags().map(|t| t.to_vec());
423 let time_in_force = primary.time_in_force();
424 let reduce_only = primary.is_reduce_only();
425 let mut primary = primary;
426 let spawned = self.spawn_market(
427 &mut primary,
428 quantity,
429 time_in_force,
430 reduce_only,
431 tags,
432 true,
433 );
434 self.submit_order(spawned.into(), None, None)?;
435
436 Ok(())
437 }
438
439 fn on_stop(&mut self) -> anyhow::Result<()> {
440 ExecutionAlgorithmNative::exec_algorithm_core_mut(self)
441 .clock_mut()
442 .cancel_timers();
443 Ok(())
444 }
445
446 fn on_reset(&mut self) -> anyhow::Result<()> {
447 self.unsubscribe_all_strategy_events();
448 ExecutionAlgorithmNative::exec_algorithm_core_mut(self).reset();
449 self.scheduled_sizes.clear();
450 Ok(())
451 }
452});
453
454fn validation_failed(detail: impl Into<String>) -> Ustr {
455 let reason = OrderDeniedReason::ValidationFailed {
456 detail: detail.into(),
457 }
458 .to_string();
459 Ustr::from(&reason)
460}
461
462#[cfg(test)]
463mod tests {
464 use std::{cell::RefCell, rc::Rc};
465
466 use indexmap::IndexMap;
467 use nautilus_common::{
468 cache::Cache,
469 clock::{Clock, TestClock},
470 component::Component,
471 enums::ComponentTrigger,
472 messages::execution::{SubmitOrder, TradingCommand},
473 msgbus::{self, MessagingSwitchboard, TypedHandler},
474 };
475 use nautilus_core::{Params, UUID4, UnixNanos};
476 use nautilus_model::{
477 enums::{OrderSide, OrderStatus, TimeInForce},
478 events::{OrderDeniedReason, OrderEventAny, order::spec::OrderCanceledSpec},
479 identifiers::{ExecAlgorithmId, InstrumentId, StrategyId, TraderId},
480 orders::{LimitOrder, MarketOrder},
481 types::Price,
482 };
483 use rstest::rstest;
484 use ustr::Ustr;
485
486 use super::*;
487
488 fn create_twap_algorithm() -> TwapAlgorithm {
489 let unique_id = format!("TWAP-{}", UUID4::new());
491 let config = TwapAlgorithmConfig {
492 exec_algorithm_id: Some(ExecAlgorithmId::new(&unique_id)),
493 ..Default::default()
494 };
495 TwapAlgorithm::new(config)
496 }
497
498 fn register_algorithm(algo: &mut TwapAlgorithm) {
499 use nautilus_common::timer::TimeEventCallback;
500
501 let trader_id = TraderId::from("TRADER-001");
502 let clock = Rc::new(RefCell::new(TestClock::new()));
503 let cache = Rc::new(RefCell::new(Cache::default()));
504
505 clock
507 .borrow_mut()
508 .register_default_handler(TimeEventCallback::Rust(std::sync::Arc::new(|_| {})));
509
510 algo.core.register(trader_id, clock, cache).unwrap();
511
512 algo.transition_state(ComponentTrigger::Initialize).unwrap();
514 algo.transition_state(ComponentTrigger::Start).unwrap();
515 algo.transition_state(ComponentTrigger::StartCompleted)
516 .unwrap();
517 }
518
519 fn add_instrument_to_cache(algo: &TwapAlgorithm) {
520 use nautilus_model::instruments::{InstrumentAny, stubs::crypto_perpetual_ethusdt};
521
522 let instrument = crypto_perpetual_ethusdt();
523 let cache_rc = algo.core.cache_rc();
524 let mut cache = cache_rc.borrow_mut();
525 cache
526 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
527 .unwrap();
528 }
529
530 fn create_market_order_with_params(params: IndexMap<Ustr, Ustr>) -> OrderAny {
531 create_market_order_with_params_and_qty(params, Quantity::from("1.0"))
532 }
533
534 fn create_market_order_with_params_and_qty(
535 params: IndexMap<Ustr, Ustr>,
536 quantity: Quantity,
537 ) -> OrderAny {
538 let client_order_id = ClientOrderId::from("O-001");
539 OrderAny::Market(MarketOrder::new(
540 TraderId::from("TRADER-001"),
541 StrategyId::from("STRAT-001"),
542 InstrumentId::from("ETHUSDT-PERP.BINANCE"),
543 client_order_id,
544 OrderSide::Buy,
545 quantity,
546 TimeInForce::Gtc,
547 UUID4::new(),
548 0.into(),
549 false,
550 false,
551 None,
552 None,
553 None,
554 None,
555 Some(ExecAlgorithmId::new("TWAP")),
556 Some(params),
557 Some(client_order_id),
558 None,
559 ))
560 }
561
562 fn assert_twap_denied(algo: &mut TwapAlgorithm, order: &OrderAny, expected_reason: &str) {
563 let strategy_id = order.strategy_id();
564 {
565 let cache_rc = algo.core.cache_rc();
566 cache_rc
567 .borrow_mut()
568 .add_order(order.clone(), None, None, false)
569 .unwrap();
570 }
571 let events = Rc::new(RefCell::new(Vec::new()));
572 let handler = TypedHandler::from({
573 let events = events.clone();
574 move |event: &OrderEventAny| events.borrow_mut().push(event.clone())
575 });
576 let topic = format!("events.order.{strategy_id}");
577 msgbus::subscribe_order_events(topic.clone().into(), handler.clone(), None);
578
579 algo.on_order(order.clone()).unwrap();
580 algo.on_order(order.clone()).unwrap();
581
582 msgbus::unsubscribe_order_events(topic.into(), &handler);
583 let cached_order = algo.cache().order(&order.client_order_id()).unwrap();
584 let events = events.borrow();
585
586 assert_eq!(cached_order.status(), OrderStatus::Denied);
587 assert_eq!(events.len(), 1);
588 assert!(matches!(
589 &events[0],
590 OrderEventAny::Denied(event)
591 if event.reason.as_str() == expected_reason
592 && event.strategy_id == strategy_id
593 && event.client_order_id == order.client_order_id()
594 ));
595 assert!(algo.scheduled_sizes.is_empty());
596 assert!(algo.clock().timer_names().is_empty());
597 }
598
599 #[rstest]
600 fn test_twap_creation() {
601 let algo = create_twap_algorithm();
602 assert!(algo.id().inner().starts_with("TWAP"));
603 assert!(algo.scheduled_sizes.is_empty());
604 }
605
606 #[rstest]
607 fn test_twap_registration() {
608 let mut algo = create_twap_algorithm();
609 register_algorithm(&mut algo);
610
611 assert_eq!(algo.trader_id(), Some(TraderId::from("TRADER-001")));
612 }
613
614 #[rstest]
615 fn test_twap_reset_clears_scheduled_sizes() {
616 let mut algo = create_twap_algorithm();
617 let primary_id = ClientOrderId::new("O-001");
618
619 algo.scheduled_sizes
620 .insert(primary_id, vec![Quantity::from("1.0")]);
621
622 assert!(!algo.scheduled_sizes.is_empty());
623
624 DataActor::on_reset(&mut algo).unwrap();
626
627 assert!(algo.scheduled_sizes.is_empty());
628 }
629
630 #[rstest]
631 fn test_twap_rejects_non_market_orders() {
632 let mut algo = create_twap_algorithm();
633 register_algorithm(&mut algo);
634
635 let order = OrderAny::Limit(LimitOrder::new(
636 TraderId::from("TRADER-001"),
637 StrategyId::from("STRAT-001"),
638 InstrumentId::from("BTC/USDT.BINANCE"),
639 ClientOrderId::from("O-001"),
640 OrderSide::Buy,
641 Quantity::from("1.0"),
642 Price::from("50000.0"),
643 TimeInForce::Gtc,
644 None, false, false, false, None, None, None, None, None, None, None, None, None, None, None, UUID4::new(),
660 0.into(),
661 ));
662
663 let reason = OrderDeniedReason::UnsupportedOrderType {
664 order_type: OrderType::Limit,
665 }
666 .to_string();
667 assert_twap_denied(&mut algo, &order, &reason);
668 }
669
670 #[rstest]
671 fn test_twap_denies_missing_instrument() {
672 let mut algo = create_twap_algorithm();
673 register_algorithm(&mut algo);
674
675 let mut params = IndexMap::new();
676 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
677 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
678 let order = create_market_order_with_params(params);
679 let reason = OrderDeniedReason::InstrumentNotFound {
680 instrument_id: order.instrument_id(),
681 }
682 .to_string();
683
684 assert_twap_denied(&mut algo, &order, &reason);
685 }
686
687 #[rstest]
688 fn test_twap_rejects_missing_params() {
689 let mut algo = create_twap_algorithm();
690 register_algorithm(&mut algo);
691
692 add_instrument_to_cache(&algo);
693
694 let order = OrderAny::Market(MarketOrder::new(
695 TraderId::from("TRADER-001"),
696 StrategyId::from("STRAT-001"),
697 InstrumentId::from("ETHUSDT-PERP.BINANCE"),
698 ClientOrderId::from("O-001"),
699 OrderSide::Buy,
700 Quantity::from("1.0"),
701 TimeInForce::Gtc,
702 UUID4::new(),
703 0.into(),
704 false,
705 false,
706 None,
707 None,
708 None,
709 None,
710 None,
711 None, None,
713 None,
714 ));
715
716 assert_twap_denied(
717 &mut algo,
718 &order,
719 "VALIDATION_FAILED: exec_algorithm_params not found",
720 );
721 }
722
723 #[rstest]
724 #[case(
725 None,
726 Some("10"),
727 "VALIDATION_FAILED: horizon_secs not found in exec_algorithm_params"
728 )]
729 #[case(
730 Some("60"),
731 None,
732 "VALIDATION_FAILED: interval_secs not found in exec_algorithm_params"
733 )]
734 #[case(
735 Some("not-a-number"),
736 Some("10"),
737 "VALIDATION_FAILED: horizon_secs=not-a-number is not a valid number"
738 )]
739 #[case(
740 Some("60"),
741 Some("not-a-number"),
742 "VALIDATION_FAILED: interval_secs=not-a-number is not a valid number"
743 )]
744 fn test_twap_denies_missing_or_malformed_schedule_parameter(
745 #[case] horizon_secs: Option<&str>,
746 #[case] interval_secs: Option<&str>,
747 #[case] expected_reason: &str,
748 ) {
749 let mut algo = create_twap_algorithm();
750 register_algorithm(&mut algo);
751 add_instrument_to_cache(&algo);
752
753 let mut params = IndexMap::new();
754
755 if let Some(horizon_secs) = horizon_secs {
756 params.insert(Ustr::from("horizon_secs"), Ustr::from(horizon_secs));
757 }
758
759 if let Some(interval_secs) = interval_secs {
760 params.insert(Ustr::from("interval_secs"), Ustr::from(interval_secs));
761 }
762
763 assert_twap_denied(
764 &mut algo,
765 &create_market_order_with_params(params),
766 expected_reason,
767 );
768 }
769
770 #[rstest]
771 fn test_twap_rejects_horizon_less_than_interval() {
772 let mut algo = create_twap_algorithm();
773 register_algorithm(&mut algo);
774
775 add_instrument_to_cache(&algo);
776
777 let mut params = IndexMap::new();
778 params.insert(Ustr::from("horizon_secs"), Ustr::from("30"));
779 params.insert(Ustr::from("interval_secs"), Ustr::from("60"));
780
781 let order = create_market_order_with_params(params);
782 assert_twap_denied(
783 &mut algo,
784 &order,
785 "VALIDATION_FAILED: horizon_secs=30 must be greater than or equal to interval_secs=60",
786 );
787 }
788
789 #[rstest]
790 fn test_twap_rejects_duplicate_order() {
791 let mut algo = create_twap_algorithm();
792 register_algorithm(&mut algo);
793
794 add_instrument_to_cache(&algo);
795
796 let mut params = IndexMap::new();
797 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
798 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
799
800 let order1 = create_market_order_with_params(params.clone());
801 let order2 = create_market_order_with_params(params);
802
803 algo.on_order(order1).unwrap();
804 let result = algo.on_order(order2);
805
806 assert!(result.is_err());
807 assert!(
808 result
809 .unwrap_err()
810 .to_string()
811 .contains("already being executed")
812 );
813 }
814
815 #[rstest]
816 fn test_twap_calculates_size_schedule_evenly() {
817 let mut algo = create_twap_algorithm();
818 register_algorithm(&mut algo);
819
820 add_instrument_to_cache(&algo);
821
822 let mut params = IndexMap::new();
824 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
825 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
826
827 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
828 let primary_id = order.client_order_id();
829
830 algo.on_order(order).unwrap();
831
832 let remaining = algo.scheduled_sizes.get(&primary_id).unwrap();
834 assert_eq!(remaining.len(), 2);
835
836 for qty in remaining {
837 assert_eq!(*qty, Quantity::from("0.4"));
838 }
839 }
840
841 #[rstest]
842 fn test_twap_reduces_cached_primary_after_first_child_spawn() {
843 let mut algo = create_twap_algorithm();
844 register_algorithm(&mut algo);
845
846 add_instrument_to_cache(&algo);
847
848 let mut params = IndexMap::new();
849 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
850 params.insert(Ustr::from("interval_secs"), Ustr::from("30"));
851
852 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
853 let primary_id = order.client_order_id();
854
855 algo.on_order(order).unwrap();
856
857 let cache = algo.cache();
858 let primary = cache.order(&primary_id).unwrap();
859 let spawned = cache.order(&ClientOrderId::from("O-001-E1")).unwrap();
860
861 assert_eq!(primary.quantity(), Quantity::from("0.6"));
862 assert_eq!(spawned.quantity(), Quantity::from("0.6"));
863 assert_eq!(spawned.exec_spawn_id(), Some(primary_id));
864 }
865
866 #[rstest]
867 fn test_twap_on_order_accepts_already_cached_primary() {
868 let mut algo = create_twap_algorithm();
869 register_algorithm(&mut algo);
870
871 add_instrument_to_cache(&algo);
872
873 let mut params = IndexMap::new();
874 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
875 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
876
877 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
878 let primary_id = order.client_order_id();
879
880 {
882 let cache_rc = algo.core.cache_rc();
883 let mut cache = cache_rc.borrow_mut();
884 cache.add_order(order.clone(), None, None, false).unwrap();
885 }
886
887 algo.on_order(order).unwrap();
888
889 assert_eq!(algo.scheduled_sizes.get(&primary_id).unwrap().len(), 2);
890 }
891
892 #[rstest]
893 fn test_twap_calculates_size_schedule_with_remainder() {
894 let mut algo = create_twap_algorithm();
895 register_algorithm(&mut algo);
896
897 add_instrument_to_cache(&algo);
898 let instrument = nautilus_model::instruments::stubs::crypto_perpetual_ethusdt();
899
900 let mut params = IndexMap::new();
902 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
903 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
904
905 let order = create_market_order_with_params(params);
906 let primary_id = order.client_order_id();
907
908 algo.on_order(order).unwrap();
909
910 let remaining = algo.scheduled_sizes.get(&primary_id).unwrap();
912 assert_eq!(remaining.len(), 3);
913 assert_eq!(
914 remaining,
915 &[
916 Quantity::from("0.333"),
917 Quantity::from("0.333"),
918 Quantity::from("0.001"),
919 ]
920 );
921
922 for quantity in remaining {
923 assert_eq!(quantity.precision, instrument.size_precision());
924 assert_eq!(quantity.raw % instrument.size_increment().raw, 0);
925 }
926
927 let first = algo
928 .cache()
929 .order(&ClientOrderId::from("O-001-E1"))
930 .unwrap()
931 .quantity();
932 let total = remaining
933 .iter()
934 .fold(first.as_decimal(), |sum, qty| sum + qty.as_decimal());
935 assert_eq!(total, Quantity::from("1.0").as_decimal());
936 }
937
938 #[rstest]
939 fn test_twap_children_use_instrument_size_precision() {
940 let mut algo = create_twap_algorithm();
941 register_algorithm(&mut algo);
942
943 add_instrument_to_cache(&algo);
944 let instrument = nautilus_model::instruments::stubs::crypto_perpetual_ethusdt();
945
946 let mut params = IndexMap::new();
947 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
948 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
949
950 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.000000000"));
951 let primary_id = order.client_order_id();
952
953 algo.on_order(order).unwrap();
954
955 let spawned = algo
956 .cache()
957 .order(&ClientOrderId::from("O-001-E1"))
958 .unwrap()
959 .quantity();
960 assert_eq!(spawned.precision, instrument.size_precision());
961 assert!(
962 algo.scheduled_sizes[&primary_id]
963 .iter()
964 .all(|quantity| quantity.precision == instrument.size_precision())
965 );
966 }
967
968 #[rstest]
969 fn test_twap_on_time_event_spawns_next_slice() {
970 let mut algo = create_twap_algorithm();
971 register_algorithm(&mut algo);
972
973 add_instrument_to_cache(&algo);
974
975 let mut params = IndexMap::new();
977 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
978 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
979
980 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
981 let primary_id = order.client_order_id();
982 let mut submit_params = Params::new();
983 submit_params.insert(
984 "routing_profile".to_string(),
985 serde_json::Value::String("intermediate-slice".to_string()),
986 );
987 algo.core
988 .remember_submit_params(primary_id, Some(submit_params.clone()));
989
990 algo.on_order(order).unwrap();
991
992 assert_eq!(algo.scheduled_sizes.get(&primary_id).unwrap().len(), 2);
994
995 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
997 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
998
999 assert_eq!(algo.scheduled_sizes.get(&primary_id).unwrap().len(), 1);
1001 assert_eq!(algo.core.submit_params(&primary_id), Some(submit_params));
1002 }
1003
1004 #[rstest]
1005 fn test_twap_data_actor_dispatch_spawns_next_slice() {
1006 let mut algo = create_twap_algorithm();
1007 register_algorithm(&mut algo);
1008
1009 add_instrument_to_cache(&algo);
1010
1011 let mut params = IndexMap::new();
1012 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1013 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1014
1015 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1016 let primary_id = order.client_order_id();
1017
1018 algo.on_order(order).unwrap();
1019 assert_eq!(algo.scheduled_sizes.get(&primary_id).unwrap().len(), 2);
1020
1021 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1023 algo.handle_time_event(&event);
1024
1025 assert_eq!(algo.scheduled_sizes.get(&primary_id).unwrap().len(), 1);
1026 }
1027
1028 #[rstest]
1029 fn test_twap_on_time_event_completes_on_final_slice() {
1030 let mut algo = create_twap_algorithm();
1031 register_algorithm(&mut algo);
1032
1033 add_instrument_to_cache(&algo);
1034
1035 let mut params = IndexMap::new();
1037 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1038 params.insert(Ustr::from("interval_secs"), Ustr::from("30"));
1039
1040 let order = create_market_order_with_params(params);
1041 let primary_id = order.client_order_id();
1042 let mut submit_params = Params::new();
1043 submit_params.insert(
1044 "routing_profile".to_string(),
1045 serde_json::Value::String("final-slice".to_string()),
1046 );
1047 algo.core
1048 .remember_submit_params(primary_id, Some(submit_params.clone()));
1049
1050 algo.on_order(order).unwrap();
1051 assert_eq!(algo.scheduled_sizes.get(&primary_id).unwrap().len(), 1);
1052 assert_eq!(
1053 algo.core.submit_params(&primary_id),
1054 Some(submit_params.clone())
1055 );
1056
1057 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1058 let handler = msgbus::TypedIntoHandler::from({
1059 let captured = received.clone();
1060 move |cmd: TradingCommand| {
1061 if let TradingCommand::SubmitOrder(cmd) = cmd {
1062 *captured.borrow_mut() = Some(cmd);
1063 }
1064 }
1065 });
1066 msgbus::register_trading_command_endpoint(
1067 MessagingSwitchboard::risk_engine_execute(),
1068 handler,
1069 );
1070
1071 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1073 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1074
1075 assert!(algo.scheduled_sizes.get(&primary_id).is_none());
1077 assert_eq!(
1078 received
1079 .borrow()
1080 .as_ref()
1081 .and_then(|cmd| cmd.params.clone()),
1082 Some(submit_params),
1083 );
1084 assert_eq!(algo.core.submit_params(&primary_id), None);
1085 }
1086
1087 #[rstest]
1088 fn test_twap_on_time_event_completes_when_primary_closed() {
1089 let mut algo = create_twap_algorithm();
1090 register_algorithm(&mut algo);
1091
1092 add_instrument_to_cache(&algo);
1093
1094 let mut params = IndexMap::new();
1095 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1096 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1097
1098 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1099 let primary_id = order.client_order_id();
1100 let mut submit_params = Params::new();
1101 submit_params.insert(
1102 "routing_profile".to_string(),
1103 serde_json::Value::String("closed-primary".to_string()),
1104 );
1105 algo.core
1106 .remember_submit_params(primary_id, Some(submit_params));
1107
1108 algo.on_order(order).unwrap();
1109 assert_eq!(algo.scheduled_sizes.get(&primary_id).unwrap().len(), 2);
1110
1111 {
1113 let cache_rc = algo.core.cache_rc();
1114 let mut cache = cache_rc.borrow_mut();
1115 let primary = cache.order(&primary_id).map(|o| o.clone()).unwrap();
1116
1117 let canceled = OrderCanceledSpec::builder()
1118 .trader_id(primary.trader_id())
1119 .strategy_id(primary.strategy_id())
1120 .instrument_id(primary.instrument_id())
1121 .client_order_id(primary.client_order_id())
1122 .build();
1123 cache
1124 .update_order(&OrderEventAny::Canceled(canceled))
1125 .unwrap();
1126 }
1127
1128 let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1130 ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1131
1132 assert!(algo.scheduled_sizes.get(&primary_id).is_none());
1134 assert_eq!(algo.core.submit_params(&primary_id), None);
1135 }
1136
1137 #[rstest]
1138 fn test_twap_on_stop_cancels_timers() {
1139 let mut algo = create_twap_algorithm();
1140 register_algorithm(&mut algo);
1141
1142 add_instrument_to_cache(&algo);
1143
1144 let mut params = IndexMap::new();
1145 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1146 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1147
1148 let order = create_market_order_with_params(params);
1149 let primary_id = order.client_order_id();
1150
1151 algo.on_order(order).unwrap();
1152
1153 assert!(
1155 algo.clock()
1156 .timer_names()
1157 .iter()
1158 .any(|name| name.as_str() == primary_id.as_str())
1159 );
1160
1161 DataActor::on_stop(&mut algo).unwrap();
1163
1164 assert!(algo.clock().timer_names().is_empty());
1166 }
1167
1168 #[rstest]
1169 fn test_twap_fractional_interval_secs() {
1170 let mut algo = create_twap_algorithm();
1171 register_algorithm(&mut algo);
1172
1173 add_instrument_to_cache(&algo);
1174
1175 let mut params = IndexMap::new();
1177 params.insert(Ustr::from("horizon_secs"), Ustr::from("3"));
1178 params.insert(Ustr::from("interval_secs"), Ustr::from("0.5"));
1179
1180 let order = create_market_order_with_params(params);
1181 let primary_id = order.client_order_id();
1182
1183 algo.on_order(order).unwrap();
1185
1186 let remaining = algo.scheduled_sizes.get(&primary_id).unwrap();
1188 assert!(remaining.len() >= 5);
1189 }
1190
1191 #[rstest]
1192 fn test_twap_submits_entire_size_when_qty_per_interval_below_size_increment() {
1193 use nautilus_model::instruments::{InstrumentAny, stubs::equity_aapl};
1194
1195 let mut algo = create_twap_algorithm();
1196 register_algorithm(&mut algo);
1197
1198 let instrument = equity_aapl();
1200 let instrument_id = instrument.id();
1201 {
1202 let cache_rc = algo.core.cache_rc();
1203 let mut cache = cache_rc.borrow_mut();
1204 cache
1205 .add_instrument(InstrumentAny::Equity(instrument))
1206 .unwrap();
1207 }
1208
1209 let mut params = IndexMap::new();
1212 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1213 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
1214
1215 let client_order_id = ClientOrderId::from("O-002");
1216 let order = OrderAny::Market(MarketOrder::new(
1217 TraderId::from("TRADER-001"),
1218 StrategyId::from("STRAT-001"),
1219 instrument_id,
1220 client_order_id,
1221 OrderSide::Buy,
1222 Quantity::from("2"),
1223 TimeInForce::Gtc,
1224 UUID4::new(),
1225 0.into(),
1226 false,
1227 false,
1228 None,
1229 None,
1230 None,
1231 None,
1232 Some(ExecAlgorithmId::new("TWAP")),
1233 Some(params),
1234 Some(client_order_id),
1235 None,
1236 ));
1237
1238 let primary_id = order.client_order_id();
1239 let mut submit_params = Params::new();
1240 submit_params.insert(
1241 "routing_profile".to_string(),
1242 serde_json::Value::String("whole-size".to_string()),
1243 );
1244 algo.core
1245 .remember_submit_params(primary_id, Some(submit_params));
1246 algo.on_order(order).unwrap();
1247
1248 assert!(algo.scheduled_sizes.get(&primary_id).is_none());
1250 assert_eq!(algo.core.submit_params(&primary_id), None);
1251 }
1252
1253 #[rstest]
1254 fn test_twap_submits_entire_size_when_qty_per_interval_below_min_quantity() {
1255 use nautilus_model::instruments::{InstrumentAny, stubs::crypto_perpetual_ethusdt};
1256
1257 let mut algo = create_twap_algorithm();
1258 register_algorithm(&mut algo);
1259
1260 let mut instrument = crypto_perpetual_ethusdt();
1261 instrument.min_quantity = Some(Quantity::from("0.5"));
1262 {
1263 let cache_rc = algo.core.cache_rc();
1264 let mut cache = cache_rc.borrow_mut();
1265 cache
1266 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
1267 .unwrap();
1268 }
1269
1270 let mut params = IndexMap::new();
1271 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1272 params.insert(Ustr::from("interval_secs"), Ustr::from("20"));
1273
1274 let order = create_market_order_with_params_and_qty(params, Quantity::from("1.2"));
1275 let primary_id = order.client_order_id();
1276 let mut submit_params = Params::new();
1277 submit_params.insert(
1278 "routing_profile".to_string(),
1279 serde_json::Value::String("minimum-quantity".to_string()),
1280 );
1281 algo.core
1282 .remember_submit_params(primary_id, Some(submit_params));
1283
1284 algo.on_order(order).unwrap();
1285
1286 assert!(algo.scheduled_sizes.get(&primary_id).is_none());
1287 assert_eq!(algo.core.submit_params(&primary_id), None);
1288 }
1289
1290 #[rstest]
1291 fn test_twap_rejects_negative_interval_secs() {
1292 let mut algo = create_twap_algorithm();
1293 register_algorithm(&mut algo);
1294
1295 add_instrument_to_cache(&algo);
1296
1297 let mut params = IndexMap::new();
1298 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1299 params.insert(Ustr::from("interval_secs"), Ustr::from("-0.5"));
1300
1301 let order = create_market_order_with_params(params);
1302 assert_twap_denied(
1303 &mut algo,
1304 &order,
1305 "VALIDATION_FAILED: interval_secs=-0.5 must be finite and positive",
1306 );
1307 }
1308
1309 #[rstest]
1310 fn test_twap_rejects_negative_horizon_secs() {
1311 let mut algo = create_twap_algorithm();
1312 register_algorithm(&mut algo);
1313
1314 add_instrument_to_cache(&algo);
1315
1316 let mut params = IndexMap::new();
1317 params.insert(Ustr::from("horizon_secs"), Ustr::from("-10"));
1318 params.insert(Ustr::from("interval_secs"), Ustr::from("1"));
1319
1320 let order = create_market_order_with_params(params);
1321 assert_twap_denied(
1322 &mut algo,
1323 &order,
1324 "VALIDATION_FAILED: horizon_secs=-10 must be finite and positive",
1325 );
1326 }
1327
1328 #[rstest]
1329 fn test_twap_rejects_zero_interval_secs() {
1330 let mut algo = create_twap_algorithm();
1331 register_algorithm(&mut algo);
1332
1333 add_instrument_to_cache(&algo);
1334
1335 let mut params = IndexMap::new();
1336 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1337 params.insert(Ustr::from("interval_secs"), Ustr::from("0"));
1338
1339 let order = create_market_order_with_params(params);
1340 assert_twap_denied(
1341 &mut algo,
1342 &order,
1343 "VALIDATION_FAILED: interval_secs=0 must be finite and positive",
1344 );
1345 }
1346
1347 #[rstest]
1348 fn test_twap_rejects_huge_finite_interval_before_submission() {
1349 let mut algo = create_twap_algorithm();
1350 register_algorithm(&mut algo);
1351
1352 add_instrument_to_cache(&algo);
1353
1354 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1355 let handler = msgbus::TypedIntoHandler::from({
1356 let captured = received.clone();
1357 move |cmd: TradingCommand| {
1358 if let TradingCommand::SubmitOrder(cmd) = cmd {
1359 *captured.borrow_mut() = Some(cmd);
1360 }
1361 }
1362 });
1363 msgbus::register_trading_command_endpoint(
1364 MessagingSwitchboard::risk_engine_execute(),
1365 handler,
1366 );
1367
1368 let mut params = IndexMap::new();
1369 params.insert(Ustr::from("horizon_secs"), Ustr::from("2e20"));
1370 params.insert(Ustr::from("interval_secs"), Ustr::from("1e20"));
1371
1372 let duration_error = Duration::try_from_secs_f64(1e20).unwrap_err();
1373 let reason = OrderDeniedReason::ValidationFailed {
1374 detail: format!(
1375 "interval_secs=100000000000000000000 is not a valid duration: {duration_error}"
1376 ),
1377 }
1378 .to_string();
1379 assert_twap_denied(&mut algo, &create_market_order_with_params(params), &reason);
1380
1381 assert!(received.borrow().is_none());
1382 }
1383
1384 #[rstest]
1385 fn test_twap_rejects_subnanosecond_interval_before_submission() {
1386 let mut algo = create_twap_algorithm();
1387 register_algorithm(&mut algo);
1388
1389 add_instrument_to_cache(&algo);
1390
1391 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1392 let handler = msgbus::TypedIntoHandler::from({
1393 let captured = received.clone();
1394 move |cmd: TradingCommand| {
1395 if let TradingCommand::SubmitOrder(cmd) = cmd {
1396 *captured.borrow_mut() = Some(cmd);
1397 }
1398 }
1399 });
1400 msgbus::register_trading_command_endpoint(
1401 MessagingSwitchboard::risk_engine_execute(),
1402 handler,
1403 );
1404
1405 let mut params = IndexMap::new();
1406 params.insert(Ustr::from("horizon_secs"), Ustr::from("2e-10"));
1407 params.insert(Ustr::from("interval_secs"), Ustr::from("1e-10"));
1408
1409 assert_twap_denied(
1410 &mut algo,
1411 &create_market_order_with_params(params),
1412 "VALIDATION_FAILED: interval_secs=0.0000000001 rounds to a zero duration",
1413 );
1414
1415 assert!(received.borrow().is_none());
1416 }
1417
1418 #[rstest]
1419 fn test_twap_rejects_interval_exceeding_timestamp_headroom_before_submission() {
1420 let mut algo = create_twap_algorithm();
1421 register_algorithm(&mut algo);
1422
1423 add_instrument_to_cache(&algo);
1424 DataActorNative::clock_mut(&mut algo)
1425 .as_any_mut()
1426 .downcast_mut::<TestClock>()
1427 .unwrap()
1428 .set_time(UnixNanos::new(u64::MAX - 500_000_000));
1429
1430 let received = Rc::new(RefCell::new(None::<SubmitOrder>));
1431 let handler = msgbus::TypedIntoHandler::from({
1432 let captured = received.clone();
1433 move |cmd: TradingCommand| {
1434 if let TradingCommand::SubmitOrder(cmd) = cmd {
1435 *captured.borrow_mut() = Some(cmd);
1436 }
1437 }
1438 });
1439 msgbus::register_trading_command_endpoint(
1440 MessagingSwitchboard::risk_engine_execute(),
1441 handler,
1442 );
1443
1444 let mut params = IndexMap::new();
1445 params.insert(Ustr::from("horizon_secs"), Ustr::from("2"));
1446 params.insert(Ustr::from("interval_secs"), Ustr::from("1"));
1447
1448 assert_twap_denied(
1449 &mut algo,
1450 &create_market_order_with_params(params),
1451 "VALIDATION_FAILED: interval_secs=1 exceeds the clock timestamp headroom",
1452 );
1453
1454 assert!(received.borrow().is_none());
1455 }
1456
1457 #[rstest]
1458 fn test_twap_rejects_nan_interval_secs() {
1459 let mut algo = create_twap_algorithm();
1460 register_algorithm(&mut algo);
1461
1462 add_instrument_to_cache(&algo);
1463
1464 let mut params = IndexMap::new();
1465 params.insert(Ustr::from("horizon_secs"), Ustr::from("60"));
1466 params.insert(Ustr::from("interval_secs"), Ustr::from("NaN"));
1467
1468 let order = create_market_order_with_params(params);
1469 assert_twap_denied(
1470 &mut algo,
1471 &order,
1472 "VALIDATION_FAILED: interval_secs=NaN must be finite and positive",
1473 );
1474 }
1475
1476 #[rstest]
1477 fn test_twap_rejects_infinity_horizon_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();
1484 params.insert(Ustr::from("horizon_secs"), Ustr::from("inf"));
1485 params.insert(Ustr::from("interval_secs"), Ustr::from("10"));
1486
1487 let order = create_market_order_with_params(params);
1488 assert_twap_denied(
1489 &mut algo,
1490 &order,
1491 "VALIDATION_FAILED: horizon_secs=inf must be finite and positive",
1492 );
1493 }
1494}