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