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