Skip to main content

nautilus_trading/algorithm/
twap.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Time-Weighted Average Price (TWAP) execution algorithm.
17//!
18//! The TWAP algorithm executes orders by evenly spreading them over a specified
19//! time horizon at regular intervals. This helps reduce market impact by avoiding
20//! concentration of trade size at any given time.
21//!
22//! # Parameters
23//!
24//! Orders submitted to this algorithm must include `exec_algorithm_params` with:
25//! - `horizon_secs`: Total execution horizon in seconds.
26//! - `interval_secs`: Interval between child orders in seconds.
27//!
28//! # Example
29//!
30//! An order with `horizon_secs=60` and `interval_secs=10` will spawn 6 child
31//! orders over 60 seconds, one every 10 seconds.
32
33use 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
56/// Configuration for [`TwapAlgorithm`].
57pub type TwapAlgorithmConfig = ExecutionAlgorithmConfig;
58
59/// Time-Weighted Average Price (TWAP) execution algorithm.
60///
61/// Executes orders by evenly spreading them over a specified time horizon,
62/// at regular intervals. The algorithm receives a primary order and spawns
63/// smaller child orders that are executed at regular intervals.
64#[derive(Debug)]
65pub struct TwapAlgorithm {
66    /// The algorithm core.
67    pub core: ExecutionAlgorithmCore,
68    /// Scheduled sizes for each primary order.
69    scheduled_sizes: AHashMap<ClientOrderId, Vec<Quantity>>,
70}
71
72impl TwapAlgorithm {
73    /// Creates a new [`TwapAlgorithm`] instance.
74    #[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    /// Completes the execution sequence for a primary order.
83    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
95// The clock and component lifecycle dispatch through the `DataActor` hooks,
96// so forward them to the `ExecutionAlgorithm` implementations.
97impl 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        // Only market orders supported
122        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        // Add primary order to cache so on_time_event can retrieve it,
330        // it is already present when routed through the engine's submit path.
331        {
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        // Single slice: submit the primary order directly
349        if is_single_slice {
350            self.submit_order(order, None, None)?;
351            self.complete_sequence(primary_id);
352            return Ok(());
353        }
354
355        // Multiple slices: spawn first child order and reduce primary
356        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        // Final slice: submit the primary order (already reduced to remaining quantity)
415        if is_final_slice {
416            self.submit_order(primary, None, None)?;
417            self.complete_sequence(primary_id);
418            return Ok(());
419        }
420
421        // Intermediate slice: spawn child order and reduce primary
422        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        // Use unique ID to avoid thread-local registry/msgbus conflicts in parallel tests
490        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        // Register a no-op default handler for timer callbacks
506        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        // Transition to Running state for tests
513        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        // Dispatch through the DataActor entry point the component lifecycle uses
625        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,  // expire_time
645            false, // post_only
646            false, // reduce_only
647            false, // quote_quantity
648            None,  // display_qty
649            None,  // emulation_trigger
650            None,  // trigger_instrument_id
651            None,  // contingency_type
652            None,  // order_list_id
653            None,  // linked_order_ids
654            None,  // parent_order_id
655            None,  // exec_algorithm_id
656            None,  // exec_algorithm_params
657            None,  // exec_spawn_id
658            None,  // tags
659            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, // No exec_algorithm_params
712            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        // 1.2 qty over 60s with 20s intervals = 3 intervals of 0.4 each (divides evenly)
823        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        // First slice spawned immediately, remaining 2 slices scheduled (no remainder)
833        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        // The engine submit path caches the primary before routing to the algorithm
881        {
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        // 1.0 qty over 60s with 20s intervals = 3 intervals
901        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        // First slice spawned, 3 remaining (2 regular + 1 remainder)
911        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        // Use qty that divides evenly: 1.2 / 3 = 0.4 each
976        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        // Verify 2 slices remain after first spawn (no remainder)
993        assert_eq!(algo.scheduled_sizes.get(&primary_id).unwrap().len(), 2);
994
995        // Simulate timer firing
996        let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
997        ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
998
999        // One slice consumed
1000        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        // Dispatch through the DataActor entry point the clock callback uses
1022        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        // 2 intervals: first spawned immediately, one in scheduled_sizes
1036        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        // Simulate timer firing for final slice
1072        let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1073        ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1074
1075        // Sequence completed, scheduled_sizes removed
1076        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        // Mark primary order as closed (canceled)
1112        {
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        // Timer fires but primary is closed
1129        let event = TimeEvent::new(primary_id.inner(), UUID4::new(), 0.into(), 0.into());
1130        ExecutionAlgorithm::on_time_event(&mut algo, &event).unwrap();
1131
1132        // Sequence should complete early since primary is closed
1133        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        // Verify timer is set
1154        assert!(
1155            algo.clock()
1156                .timer_names()
1157                .iter()
1158                .any(|name| name.as_str() == primary_id.as_str())
1159        );
1160
1161        // Stop through the DataActor entry point the component lifecycle uses
1162        DataActor::on_stop(&mut algo).unwrap();
1163
1164        // Timer should be canceled
1165        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        // Use fractional interval like Python tests: 3 second horizon, 0.5 second interval
1176        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        // Should not error - fractional seconds should parse correctly
1184        algo.on_order(order).unwrap();
1185
1186        // 3 / 0.5 = 6 intervals, first spawned immediately, 5 remaining (plus possible remainder)
1187        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        // Use equity with size_increment of 1 (whole shares only)
1199        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        // 2 shares over 60s with 10s intervals = 6 intervals
1210        // 2 / 6 = 0.333... which is less than size_increment of 1
1211        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        // Should submit entire size directly (no scheduling)
1249        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}