1use std::num::NonZeroUsize;
17
18use ahash::AHashSet;
19use nautilus_common::{
20 actor::{DataActor, DataActorNative},
21 config::ConfigError,
22 enums::LogColor,
23 log_info, log_warn,
24};
25use nautilus_core::UnixNanos;
26use nautilus_model::{
27 data::{Bar, IndexPriceUpdate, MarkPriceUpdate, OrderBookDeltas, QuoteTick, TradeTick},
28 enums::{ContingencyType, OrderSide, OrderStatus, OrderType, TimeInForce},
29 events::{
30 OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied, OrderEvent, OrderEventAny,
31 OrderExpired, OrderFilled, OrderModifyRejected, OrderRejected, OrderUpdated,
32 },
33 identifiers::{ClientId, ClientOrderId, InstrumentId, StrategyId},
34 instruments::{Instrument, InstrumentAny},
35 orderbook::OrderBook,
36 orders::{Order, OrderAny, OrderCore},
37 types::{Price, Quantity},
38};
39use nautilus_trading::{
40 nautilus_strategy,
41 strategy::{Strategy, StrategyCore, StrategyNative},
42};
43use rust_decimal::{Decimal, prelude::ToPrimitive};
44
45use super::config::ExecTesterConfig;
46use crate::testers::timestamps::warn_if_implausible_unix_nanos;
47
48#[derive(Debug)]
57#[expect(
58 clippy::struct_excessive_bools,
59 reason = "tester state tracks independent execution scenarios"
60)]
61pub struct ExecTester {
62 pub(super) core: StrategyCore,
63 pub(super) config: ExecTesterConfig,
64 pub(super) instrument: Option<InstrumentAny>,
65 pub(super) price_offset: Option<u64>,
66 pub(super) preinitialized_market_data: bool,
67
68 pub(super) buy_order: Option<OrderAny>,
70 pub(super) sell_order: Option<OrderAny>,
71 pub(super) buy_stop_order: Option<OrderAny>,
72 pub(super) sell_stop_order: Option<OrderAny>,
73 pub(super) open_position_submitted: bool,
74
75 pub(super) modify_rejected_attempted: bool,
78 pub(super) pending_open_position_qty: Option<Decimal>,
79 pub(super) buy_stop_cancel_replace_attempted: bool,
80 pub(super) sell_stop_cancel_replace_attempted: bool,
81 pub(super) buy_limit_maintenance_state: LimitOrderMaintenanceState,
82 pub(super) sell_limit_maintenance_state: LimitOrderMaintenanceState,
83}
84
85#[derive(Debug, Clone, Copy, PartialEq, Eq)]
86pub(super) enum LimitOrderMaintenanceState {
87 Disabled,
88 Ready,
89 AwaitingAcceptance {
90 client_order_id: ClientOrderId,
91 },
92 Active {
93 client_order_id: ClientOrderId,
94 },
95 AwaitingUpdate {
96 client_order_id: ClientOrderId,
97 price: Price,
98 },
99 AwaitingCancel {
100 client_order_id: ClientOrderId,
101 replacement_price: Price,
102 },
103 AwaitingReplacementAcceptance {
104 client_order_id: ClientOrderId,
105 price: Price,
106 },
107 Completed,
108 Failed,
109}
110
111nautilus_strategy!(ExecTester, {
112 fn external_order_claims(&self) -> Option<Vec<InstrumentId>> {
113 self.config.base.external_order_claims.clone()
114 }
115
116 fn on_order_accepted(&mut self, event: OrderAccepted) {
117 self.handle_limit_order_accepted(event.client_order_id);
118 }
119
120 fn on_order_updated(&mut self, event: OrderUpdated) {
121 self.handle_limit_order_updated(event.client_order_id, event.price);
122 }
123
124 fn on_order_canceled(&mut self, event: &OrderCanceled) {
125 self.handle_limit_order_canceled(event.client_order_id);
126 }
127
128 fn on_order_modify_rejected(&mut self, event: OrderModifyRejected) {
129 self.finish_limit_order_maintenance(event.client_order_id);
130 }
131
132 fn on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {
133 self.finish_limit_order_maintenance(event.client_order_id);
134 }
135
136 fn on_order_denied(&mut self, event: OrderDenied) {
137 self.handle_limit_order_closed(event.client_order_id);
138 }
139
140 fn on_order_rejected(&mut self, event: OrderRejected) {
141 self.handle_limit_order_closed(event.client_order_id);
142 }
143
144 fn on_order_expired(&mut self, event: OrderExpired) {
145 self.handle_limit_order_closed(event.client_order_id);
146 }
147
148 fn on_order_filled(&mut self, event: &OrderFilled) {
149 let is_closed = self
150 .cache()
151 .order(&event.client_order_id)
152 .is_none_or(|order| order.is_closed());
153 if is_closed {
154 self.handle_limit_order_closed(event.client_order_id);
155 }
156 }
157
158 fn on_order_event(&mut self, event: OrderEventAny) {
159 let event = event.into_boxed();
160 warn_if_implausible_unix_nanos(
161 "order event",
162 OrderEvent::ts_event(&*event),
163 OrderEvent::ts_init(&*event),
164 );
165 }
166});
167
168impl DataActor for ExecTester {
169 fn on_start(&mut self) -> anyhow::Result<()> {
170 Strategy::on_start(self)?;
171
172 let instrument_id = self.config.instrument_id;
173 let client_id = self.config.client_id;
174
175 let instrument = self.cache().instrument(&instrument_id);
176
177 if let Some(inst) = instrument {
178 self.initialize_with_instrument(inst, true)?;
179 } else {
180 log::info!("Instrument {instrument_id} not in cache, subscribing...");
181 self.subscribe_instrument(instrument_id, client_id, None);
182
183 if self.config.subscribe_quotes {
186 self.subscribe_quotes(instrument_id, client_id, None);
187 }
188
189 if self.config.subscribe_trades {
190 self.subscribe_trades(instrument_id, client_id, None);
191 }
192 self.preinitialized_market_data =
193 self.config.subscribe_quotes || self.config.subscribe_trades;
194 }
195
196 Ok(())
197 }
198
199 fn on_instrument(&mut self, instrument: &InstrumentAny) -> anyhow::Result<()> {
200 warn_if_implausible_unix_nanos("instrument", instrument.ts_event(), instrument.ts_init());
201
202 if instrument.id() == self.config.instrument_id && self.instrument.is_none() {
203 let id = instrument.id();
204 log::info!("Received instrument {id}, initializing...");
205 self.initialize_with_instrument(instrument.clone(), !self.preinitialized_market_data)?;
206 }
207 Ok(())
208 }
209
210 fn on_stop(&mut self) -> anyhow::Result<()> {
211 if self.config.dry_run {
212 log_warn!("Dry run mode, skipping cancel all orders and close all positions");
213 return Ok(());
214 }
215
216 let instrument_id = self.config.instrument_id;
217 let client_id = self.config.client_id;
218 let strategy_id = StrategyId::from(
219 DataActorNative::core(&self.core)
220 .actor_id()
221 .inner()
222 .as_str(),
223 );
224
225 if self.config.cancel_orders_on_stop {
226 self.cancel_active_orders(instrument_id, strategy_id, client_id);
227 }
228
229 if self.config.close_positions_on_stop {
230 let time_in_force = self
231 .config
232 .close_positions_time_in_force
233 .or(Some(TimeInForce::Gtc));
234
235 if let Err(e) = self.close_positions_on_stop(strategy_id, time_in_force) {
236 log::error!("Failed to close all positions: {e}");
237 }
238 }
239
240 if self.config.can_unsubscribe && self.instrument.is_some() {
241 if self.config.subscribe_quotes {
242 self.unsubscribe_quotes(instrument_id, client_id, None);
243 }
244
245 if self.config.subscribe_trades {
246 self.unsubscribe_trades(instrument_id, client_id, None);
247 }
248
249 if self.config.subscribe_book {
250 self.unsubscribe_book_at_interval(
251 instrument_id,
252 NonZeroUsize::new(self.config.book_interval_ms).ok_or_else(|| {
253 ConfigError::range("book_interval_ms", "must be positive, was 0")
254 })?,
255 client_id,
256 None,
257 );
258 }
259 }
260
261 Ok(())
262 }
263
264 fn on_quote(&mut self, quote: &QuoteTick) -> anyhow::Result<()> {
265 warn_if_implausible_unix_nanos("quote", quote.ts_event, quote.ts_init);
266
267 if self.config.log_data {
268 log_info!("{quote:?}", color = LogColor::Cyan);
269 }
270
271 if quote.instrument_id == self.config.instrument_id
272 && self.config.open_position_on_first_quote
273 {
274 self.submit_pending_open_position();
275 }
276
277 self.maintain_orders(quote.bid_price, quote.ask_price);
278 Ok(())
279 }
280
281 fn on_trade(&mut self, trade: &TradeTick) -> anyhow::Result<()> {
282 warn_if_implausible_unix_nanos("trade", trade.ts_event, trade.ts_init);
283
284 if self.config.log_data {
285 log_info!("{trade:?}", color = LogColor::Cyan);
286 }
287 Ok(())
288 }
289
290 fn on_book(&mut self, book: &OrderBook) -> anyhow::Result<()> {
291 if self.config.log_data {
292 let num_levels = self.config.book_levels_to_print;
293 let instrument_id = book.instrument_id;
294 let book_str = book.pprint(num_levels, None);
295 log_info!("\n{instrument_id}\n{book_str}", color = LogColor::Cyan);
296
297 if DataActorNative::core(&self.core).is_registered() {
299 let cache = self.cache();
300 if let Some(own_book) = cache.own_order_book(&instrument_id) {
301 let own_book_str = own_book.pprint(num_levels, None);
302 log_info!(
303 "\n{instrument_id} (own)\n{own_book_str}",
304 color = LogColor::Magenta
305 );
306 }
307 }
308 }
309
310 let Some(best_bid) = book.best_bid_price() else {
311 return Ok(()); };
313 let Some(best_ask) = book.best_ask_price() else {
314 return Ok(()); };
316
317 self.maintain_orders(best_bid, best_ask);
318 Ok(())
319 }
320
321 fn on_book_deltas(&mut self, deltas: &OrderBookDeltas) -> anyhow::Result<()> {
322 warn_if_implausible_unix_nanos("book deltas", deltas.ts_event, deltas.ts_init);
323
324 if self.config.log_data {
325 log_info!("{deltas:?}", color = LogColor::Cyan);
326 }
327 Ok(())
328 }
329
330 fn on_bar(&mut self, bar: &Bar) -> anyhow::Result<()> {
331 warn_if_implausible_unix_nanos("bar", bar.ts_event, bar.ts_init);
332
333 if self.config.log_data {
334 log_info!("{bar:?}", color = LogColor::Cyan);
335 }
336 Ok(())
337 }
338
339 fn on_mark_price(&mut self, mark_price: &MarkPriceUpdate) -> anyhow::Result<()> {
340 warn_if_implausible_unix_nanos("mark price", mark_price.ts_event, mark_price.ts_init);
341
342 if self.config.log_data {
343 log_info!("{mark_price:?}", color = LogColor::Cyan);
344 }
345 Ok(())
346 }
347
348 fn on_index_price(&mut self, index_price: &IndexPriceUpdate) -> anyhow::Result<()> {
349 warn_if_implausible_unix_nanos("index price", index_price.ts_event, index_price.ts_init);
350
351 if self.config.log_data {
352 log_info!("{index_price:?}", color = LogColor::Cyan);
353 }
354 Ok(())
355 }
356}
357
358impl ExecTester {
359 #[must_use]
361 pub fn new(config: ExecTesterConfig) -> Self {
362 let pending_open_position_qty = config.open_position_on_start_qty;
363 let buy_limit_maintenance_state = LimitOrderMaintenanceState::new(
364 config.enable_limit_buys
365 && (config.trigger_limit_order_maintenance_once
366 || config.cancel_replace_orders_to_maintain_tob_offset),
367 );
368 let sell_limit_maintenance_state = LimitOrderMaintenanceState::new(
369 config.enable_limit_sells
370 && (config.trigger_limit_order_maintenance_once
371 || config.cancel_replace_orders_to_maintain_tob_offset),
372 );
373
374 Self {
375 core: StrategyCore::new(config.base.clone()),
376 config,
377 instrument: None,
378 price_offset: None,
379 preinitialized_market_data: false,
380 buy_order: None,
381 sell_order: None,
382 buy_stop_order: None,
383 sell_stop_order: None,
384 open_position_submitted: false,
385 modify_rejected_attempted: false,
386 pending_open_position_qty,
387 buy_stop_cancel_replace_attempted: false,
388 sell_stop_cancel_replace_attempted: false,
389 buy_limit_maintenance_state,
390 sell_limit_maintenance_state,
391 }
392 }
393
394 fn initialize_with_instrument(
395 &mut self,
396 instrument: InstrumentAny,
397 subscribe_market_data: bool,
398 ) -> anyhow::Result<()> {
399 let instrument_id = self.config.instrument_id;
400 let client_id = self.config.client_id;
401
402 self.price_offset = Some(self.get_price_offset(&instrument));
403 self.instrument = Some(instrument);
404
405 if subscribe_market_data && self.config.subscribe_quotes {
406 self.subscribe_quotes(instrument_id, client_id, None);
407 }
408
409 if subscribe_market_data && self.config.subscribe_trades {
410 self.subscribe_trades(instrument_id, client_id, None);
411 }
412
413 if self.config.subscribe_book {
414 self.subscribe_book_at_interval(
415 instrument_id,
416 self.config.book_type,
417 self.config
418 .book_depth
419 .map(|depth| {
420 NonZeroUsize::new(depth).ok_or_else(|| {
421 ConfigError::range("book_depth", "must be positive, was 0")
422 })
423 })
424 .transpose()?,
425 NonZeroUsize::new(self.config.book_interval_ms).ok_or_else(|| {
426 ConfigError::range("book_interval_ms", "must be positive, was 0")
427 })?,
428 client_id,
429 None,
430 );
431 }
432
433 if let Some(qty) = self.pending_open_position_qty {
434 let quote_ready = {
435 let cache = self.cache();
436 cache.quote(&instrument_id).is_some()
437 };
438
439 if self.config.open_position_on_first_quote
440 && self.config.subscribe_quotes
441 && !quote_ready
442 {
443 log::info!("Waiting for first quote before opening {instrument_id} position");
444 } else {
445 self.pending_open_position_qty = None;
446 self.open_position(qty)?;
447 self.open_position_submitted = true;
448 }
449 }
450
451 Ok(())
452 }
453
454 pub(super) fn get_price_offset(&self, _instrument: &InstrumentAny) -> u64 {
455 self.config.tob_offset_ticks
456 }
457
458 fn expire_time_from_delta(&self, mins: u64) -> UnixNanos {
459 let current_ns = DataActorNative::core(&self.core).timestamp_ns();
460 let delta_ns = mins.saturating_mul(60).saturating_mul(1_000_000_000);
461 UnixNanos::from(current_ns.as_u64() + delta_ns)
462 }
463
464 fn resolve_time_in_force(
465 &self,
466 tif_override: Option<TimeInForce>,
467 ) -> (TimeInForce, Option<UnixNanos>) {
468 match (tif_override, self.config.order_expire_time_delta_mins) {
469 (Some(TimeInForce::Gtd), Some(mins)) => {
470 (TimeInForce::Gtd, Some(self.expire_time_from_delta(mins)))
471 }
472 (Some(TimeInForce::Gtd), None) => {
473 log_warn!(
474 "GTD time in force requires order_expire_time_delta_mins, falling back to GTC"
475 );
476 (TimeInForce::Gtc, None)
477 }
478 (Some(tif), _) => (tif, None),
479 (None, Some(mins)) => (TimeInForce::Gtd, Some(self.expire_time_from_delta(mins))),
480 (None, None) => (TimeInForce::Gtc, None),
481 }
482 }
483
484 fn submit_pending_open_position(&mut self) {
485 if self.instrument.is_none() {
486 return;
487 }
488
489 let Some(qty) = self.pending_open_position_qty.take() else {
490 return;
491 };
492
493 if let Err(e) = self.open_position(qty) {
494 log::error!("Failed to submit pending open position: {e}");
495 } else {
496 self.open_position_submitted = true;
497 }
498 }
499
500 pub(super) fn is_order_active(order: &OrderAny) -> bool {
501 order.is_active_local() || order.is_inflight() || order.is_open()
502 }
503
504 pub(super) fn limit_order_is_one_shot(&self) -> bool {
505 self.config.test_reject_post_only
506 || self.config.limit_aggressive
507 || self.config.order_expire_time_delta_mins.is_some()
508 || matches!(
509 self.config.limit_time_in_force,
510 Some(TimeInForce::Ioc | TimeInForce::Fok)
511 )
512 }
513
514 pub(super) fn stop_order_is_one_shot(&self) -> bool {
515 self.config.order_expire_time_delta_mins.is_some()
516 || matches!(
517 self.config.stop_time_in_force,
518 Some(TimeInForce::Ioc | TimeInForce::Fok)
519 )
520 || matches!(self.config.stop_order_type, OrderType::TrailingStopMarket)
521 }
522
523 pub(super) fn get_order_trigger_price(order: &OrderAny) -> Option<Price> {
524 order.trigger_price()
525 }
526
527 fn modify_stop_order(
528 &mut self,
529 order: &OrderAny,
530 trigger_price: Price,
531 limit_price: Option<Price>,
532 ) -> anyhow::Result<()> {
533 let client_id = self.config.client_id;
534
535 match order {
536 OrderAny::StopMarket(_)
537 | OrderAny::MarketIfTouched(_)
538 | OrderAny::TrailingStopMarket(_) => self.modify_order(
539 order.client_order_id(),
540 None,
541 None,
542 Some(trigger_price),
543 client_id,
544 None,
545 ),
546 OrderAny::StopLimit(_) | OrderAny::LimitIfTouched(_) => self.modify_order(
547 order.client_order_id(),
548 None,
549 limit_price,
550 Some(trigger_price),
551 client_id,
552 None,
553 ),
554 _ => {
555 log_warn!("Cannot modify order of type {:?}", order.order_type());
556 Ok(())
557 }
558 }
559 }
560
561 fn submit_order_apply_params(&mut self, order: OrderAny) -> anyhow::Result<()> {
563 let client_id = self.config.client_id;
564 if let Some(params) = &self.config.order_params {
565 self.submit_order(order, None, client_id, Some(params.clone()))
566 } else {
567 self.submit_order(order, None, client_id, None)
568 }
569 }
570
571 pub(super) fn maintain_orders(&mut self, best_bid: Price, best_ask: Price) {
573 if self.instrument.is_none() || self.config.dry_run {
574 return;
575 }
576
577 if self.config.batch_submit_limit_pair
578 && self.config.enable_limit_buys
579 && self.config.enable_limit_sells
580 {
581 self.maintain_batch_limit_pair(best_bid, best_ask);
582 return;
583 }
584
585 if self.config.enable_limit_buys {
586 self.maintain_buy_orders(best_bid, best_ask);
587 }
588
589 if self.config.enable_limit_sells {
590 self.maintain_sell_orders(best_bid, best_ask);
591 }
592
593 if self.config.enable_stop_buys {
594 self.maintain_stop_buy_orders(best_bid, best_ask);
595 }
596
597 if self.config.enable_stop_sells {
598 self.maintain_stop_sell_orders(best_bid, best_ask);
599 }
600 }
601
602 fn refresh_tracked_order(&mut self, side: OrderSide) {
606 let cid = match side {
607 OrderSide::Buy => self.buy_order.as_ref().map(OrderAny::client_order_id),
608 OrderSide::Sell => self.sell_order.as_ref().map(OrderAny::client_order_id),
609 };
610 let Some(cid) = cid else {
611 return;
612 };
613 let latest = self.cache().order(&cid);
614 if let Some(latest) = latest {
615 match side {
616 OrderSide::Buy => self.buy_order = Some(latest),
617 OrderSide::Sell => self.sell_order = Some(latest),
618 }
619 }
620 }
621
622 fn refresh_tracked_stop_order(&mut self, side: OrderSide) {
623 let cid = match side {
624 OrderSide::Buy => self.buy_stop_order.as_ref().map(OrderAny::client_order_id),
625 OrderSide::Sell => self.sell_stop_order.as_ref().map(OrderAny::client_order_id),
626 };
627 let Some(cid) = cid else {
628 return;
629 };
630 let latest = self.cache().order(&cid);
631 if let Some(latest) = latest {
632 match side {
633 OrderSide::Buy => self.buy_stop_order = Some(latest),
634 OrderSide::Sell => self.sell_stop_order = Some(latest),
635 }
636 }
637 }
638
639 fn maintain_buy_orders(&mut self, best_bid: Price, best_ask: Price) {
641 self.refresh_tracked_order(OrderSide::Buy);
645
646 let Some(instrument) = &self.instrument else {
647 return;
648 };
649 let Some(price_offset_ticks) = self.price_offset else {
650 return;
651 };
652
653 let increment = instrument.price_increment();
654 let precision = instrument.price_precision();
655
656 let cross_spread = self.config.test_reject_post_only || self.config.limit_aggressive;
661 let unclamped_price = if cross_spread {
662 add_price_ticks(best_ask, increment, price_offset_ticks, precision)
663 } else {
664 sub_price_ticks(best_bid, increment, price_offset_ticks, precision)
665 };
666 let price = clamp_price_to_range(
667 unclamped_price,
668 instrument,
669 self.config.clamp_to_instrument_price_range,
670 );
671
672 if self.limit_order_maintenance_suppressed(OrderSide::Buy) {
673 return;
674 }
675
676 let needs_new_order = match &self.buy_order {
677 None => true,
678 Some(order) => !Self::is_order_active(order) && !self.limit_order_is_one_shot(),
679 };
680
681 if needs_new_order {
682 let result = if self.config.enable_brackets {
683 self.submit_bracket_order(OrderSide::Buy, price)
684 } else {
685 self.submit_limit_order(OrderSide::Buy, price)
686 };
687
688 if let Err(e) = result {
689 log::error!("Failed to submit buy order: {e}");
690 }
691 } else if let Some(order) = &self.buy_order
692 && order.venue_order_id().is_some()
693 && !order.is_pending_update()
694 && !order.is_pending_cancel()
695 {
696 let client_id = self.config.client_id;
697
698 if self.config.test_modify_rejected && !self.modify_rejected_attempted {
701 self.modify_rejected_attempted = true;
702 let order_clone = order.clone();
703 let bumped = clamp_price_to_range(
704 add_price_ticks(price, increment, 1, precision),
705 instrument,
706 self.config.clamp_to_instrument_price_range,
707 );
708
709 if let Err(e) = self.modify_order(
710 order_clone.client_order_id(),
711 None,
712 Some(bumped),
713 None,
714 client_id,
715 None,
716 ) {
717 log::error!("Failed to submit test modify on buy order: {e}");
718 }
719 return;
720 }
721
722 if let Some(order_price) = order.price()
723 && order_price < price
724 {
725 if self.config.modify_orders_to_maintain_tob_offset {
726 let order_clone = order.clone();
727 if let Err(e) = self.modify_order(
728 order_clone.client_order_id(),
729 None,
730 Some(price),
731 None,
732 client_id,
733 None,
734 ) {
735 log::error!("Failed to modify buy order: {e}");
736 }
737 } else if self.config.cancel_replace_orders_to_maintain_tob_offset {
738 let client_order_id = order.client_order_id();
739 self.start_limit_order_cancel_replace(OrderSide::Buy, client_order_id, price);
740 }
741 }
742 }
743 }
744
745 fn maintain_sell_orders(&mut self, best_bid: Price, best_ask: Price) {
747 self.refresh_tracked_order(OrderSide::Sell);
750
751 let Some(instrument) = &self.instrument else {
752 return;
753 };
754 let Some(price_offset_ticks) = self.price_offset else {
755 return;
756 };
757
758 let increment = instrument.price_increment();
759 let precision = instrument.price_precision();
760
761 let cross_spread = self.config.test_reject_post_only || self.config.limit_aggressive;
763 let unclamped_price = if cross_spread {
764 sub_price_ticks(best_bid, increment, price_offset_ticks, precision)
765 } else {
766 add_price_ticks(best_ask, increment, price_offset_ticks, precision)
767 };
768 let price = clamp_price_to_range(
769 unclamped_price,
770 instrument,
771 self.config.clamp_to_instrument_price_range,
772 );
773
774 if self.limit_order_maintenance_suppressed(OrderSide::Sell) {
775 return;
776 }
777
778 let needs_new_order = match &self.sell_order {
779 None => true,
780 Some(order) => !Self::is_order_active(order) && !self.limit_order_is_one_shot(),
781 };
782
783 if needs_new_order {
784 let result = if self.config.enable_brackets {
785 self.submit_bracket_order(OrderSide::Sell, price)
786 } else {
787 self.submit_limit_order(OrderSide::Sell, price)
788 };
789
790 if let Err(e) = result {
791 log::error!("Failed to submit sell order: {e}");
792 }
793 } else if let Some(order) = &self.sell_order
794 && order.venue_order_id().is_some()
795 && !order.is_pending_update()
796 && !order.is_pending_cancel()
797 {
798 let client_id = self.config.client_id;
799
800 if self.config.test_modify_rejected && !self.modify_rejected_attempted {
802 self.modify_rejected_attempted = true;
803 let order_clone = order.clone();
804 let bumped = clamp_price_to_range(
805 sub_price_ticks(price, increment, 1, precision),
806 instrument,
807 self.config.clamp_to_instrument_price_range,
808 );
809
810 if let Err(e) = self.modify_order(
811 order_clone.client_order_id(),
812 None,
813 Some(bumped),
814 None,
815 client_id,
816 None,
817 ) {
818 log::error!("Failed to submit test modify on sell order: {e}");
819 }
820 return;
821 }
822
823 if let Some(order_price) = order.price()
824 && order_price > price
825 {
826 if self.config.modify_orders_to_maintain_tob_offset {
827 let order_clone = order.clone();
828 if let Err(e) = self.modify_order(
829 order_clone.client_order_id(),
830 None,
831 Some(price),
832 None,
833 client_id,
834 None,
835 ) {
836 log::error!("Failed to modify sell order: {e}");
837 }
838 } else if self.config.cancel_replace_orders_to_maintain_tob_offset {
839 let client_order_id = order.client_order_id();
840 self.start_limit_order_cancel_replace(OrderSide::Sell, client_order_id, price);
841 }
842 }
843 }
844 }
845
846 fn handle_limit_order_accepted(&mut self, client_order_id: ClientOrderId) {
847 let Some(side) = self.limit_order_side_for_state(client_order_id) else {
848 return;
849 };
850
851 match self.limit_order_maintenance_state(side) {
852 LimitOrderMaintenanceState::AwaitingAcceptance { .. } => {
853 self.refresh_tracked_order(side);
854 if self.config.trigger_limit_order_maintenance_once {
855 self.trigger_limit_order_maintenance(side, client_order_id);
856 } else {
857 self.set_limit_order_maintenance_state(
858 side,
859 LimitOrderMaintenanceState::Active { client_order_id },
860 );
861 }
862 }
863 LimitOrderMaintenanceState::AwaitingReplacementAcceptance { price, .. } => {
864 self.refresh_tracked_order(side);
865 let accepted_at_expected_price = self
866 .tracked_limit_order(side)
867 .is_some_and(|order| order.price() == Some(price) && order.is_open());
868 if !accepted_at_expected_price {
869 log::error!(
870 "Replacement {side} order {client_order_id} was accepted with unexpected state"
871 );
872 self.set_limit_order_maintenance_state(
873 side,
874 LimitOrderMaintenanceState::Failed,
875 );
876 return;
877 }
878 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Completed);
879 }
880 _ => {}
881 }
882 }
883
884 fn trigger_limit_order_maintenance(&mut self, side: OrderSide, client_order_id: ClientOrderId) {
885 let Some(order_price) = self.tracked_limit_order(side).and_then(Order::price) else {
886 log::error!("Accepted {side} limit order {client_order_id} has no price");
887 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
888 return;
889 };
890 let Some(price) = self.limit_order_trigger_price(side, order_price) else {
891 log::error!(
892 "Cannot find a different valid price for accepted {side} limit order {client_order_id}"
893 );
894 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
895 return;
896 };
897
898 if self.config.modify_orders_to_maintain_tob_offset {
899 self.set_limit_order_maintenance_state(
900 side,
901 LimitOrderMaintenanceState::AwaitingUpdate {
902 client_order_id,
903 price,
904 },
905 );
906
907 if let Err(e) = self.modify_order(
908 client_order_id,
909 None,
910 Some(price),
911 None,
912 self.config.client_id,
913 None,
914 ) {
915 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
916 log::error!("Failed to submit one-shot {side} limit order modify: {e}");
917 }
918 } else if self.config.cancel_replace_orders_to_maintain_tob_offset {
919 self.start_limit_order_cancel_replace(side, client_order_id, price);
920 } else {
921 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
922 log::error!("No limit order maintenance mode configured for one-shot trigger");
923 }
924 }
925
926 fn start_limit_order_cancel_replace(
927 &mut self,
928 side: OrderSide,
929 client_order_id: ClientOrderId,
930 replacement_price: Price,
931 ) {
932 let state = self.limit_order_maintenance_state(side);
933 if !matches!(
934 state,
935 LimitOrderMaintenanceState::AwaitingAcceptance {
936 client_order_id: expected,
937 } | LimitOrderMaintenanceState::Active {
938 client_order_id: expected,
939 } if expected == client_order_id
940 ) {
941 return;
942 }
943
944 self.set_limit_order_maintenance_state(
945 side,
946 LimitOrderMaintenanceState::AwaitingCancel {
947 client_order_id,
948 replacement_price,
949 },
950 );
951
952 if let Err(e) = self.cancel_order(client_order_id, self.config.client_id, None) {
953 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
954 log::error!("Failed to cancel {side} limit order for replacement: {e}");
955 }
956 }
957
958 fn handle_limit_order_updated(
959 &mut self,
960 client_order_id: ClientOrderId,
961 updated_price: Option<Price>,
962 ) {
963 let Some(side) = self.limit_order_side_for_state(client_order_id) else {
964 return;
965 };
966 let LimitOrderMaintenanceState::AwaitingUpdate { price, .. } =
967 self.limit_order_maintenance_state(side)
968 else {
969 return;
970 };
971
972 self.refresh_tracked_order(side);
973 let updated_at_expected_price = updated_price == Some(price)
974 && self
975 .tracked_limit_order(side)
976 .is_some_and(|order| order.price() == Some(price) && order.is_open());
977 if !updated_at_expected_price {
978 log::error!(
979 "Updated {side} limit order {client_order_id} has unexpected price or state"
980 );
981 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
982 return;
983 }
984 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Completed);
985 }
986
987 fn handle_limit_order_canceled(&mut self, client_order_id: ClientOrderId) {
988 let Some(side) = self.limit_order_side_for_state(client_order_id) else {
989 return;
990 };
991 let LimitOrderMaintenanceState::AwaitingCancel {
992 replacement_price, ..
993 } = self.limit_order_maintenance_state(side)
994 else {
995 self.handle_limit_order_closed(client_order_id);
996 return;
997 };
998
999 self.refresh_tracked_order(side);
1000 let Some(leaves_qty) = self
1001 .tracked_limit_order(side)
1002 .filter(|order| order.client_order_id() == client_order_id)
1003 .map(Order::leaves_qty)
1004 else {
1005 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
1006 log::error!("Canceled {side} limit order {client_order_id} is not tracked");
1007 return;
1008 };
1009
1010 if leaves_qty.is_zero() {
1011 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
1012 log::warn!(
1013 "Canceled {side} limit order {client_order_id} has no remaining quantity to replace"
1014 );
1015 return;
1016 }
1017
1018 if let Err(e) = self.submit_limit_order_replacement(side, replacement_price, leaves_qty) {
1019 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
1020 log::error!("Failed to submit replacement {side} limit order: {e}");
1021 }
1022 }
1023
1024 fn finish_limit_order_maintenance(&mut self, client_order_id: ClientOrderId) {
1025 let Some(side) = self.limit_order_side_for_state(client_order_id) else {
1026 return;
1027 };
1028
1029 if matches!(
1030 self.limit_order_maintenance_state(side),
1031 LimitOrderMaintenanceState::AwaitingUpdate { .. }
1032 | LimitOrderMaintenanceState::AwaitingCancel { .. }
1033 | LimitOrderMaintenanceState::AwaitingReplacementAcceptance { .. }
1034 ) {
1035 self.set_limit_order_maintenance_state(side, LimitOrderMaintenanceState::Failed);
1036 }
1037 }
1038
1039 fn handle_limit_order_closed(&mut self, client_order_id: ClientOrderId) {
1040 let Some(side) = self.limit_order_side_for_state(client_order_id) else {
1041 return;
1042 };
1043 self.refresh_tracked_order(side);
1044 let state = match self.limit_order_maintenance_state(side) {
1045 LimitOrderMaintenanceState::AwaitingAcceptance { .. }
1046 | LimitOrderMaintenanceState::Active { .. } => LimitOrderMaintenanceState::Ready,
1047 LimitOrderMaintenanceState::AwaitingUpdate { .. }
1048 | LimitOrderMaintenanceState::AwaitingCancel { .. }
1049 | LimitOrderMaintenanceState::AwaitingReplacementAcceptance { .. } => {
1050 LimitOrderMaintenanceState::Failed
1051 }
1052 state => state,
1053 };
1054 self.set_limit_order_maintenance_state(side, state);
1055 }
1056
1057 fn limit_order_trigger_price(&self, side: OrderSide, order_price: Price) -> Option<Price> {
1058 let instrument = self.instrument.as_ref()?;
1059 let increment = instrument.price_increment();
1060 let precision = instrument.price_precision();
1061 let passive = match side {
1062 OrderSide::Buy => sub_price_ticks(order_price, increment, 1, precision),
1063 OrderSide::Sell => add_price_ticks(order_price, increment, 1, precision),
1064 };
1065 let passive = clamp_price_to_range(passive, instrument, true);
1066 if passive != order_price {
1067 return Some(passive);
1068 }
1069
1070 let aggressive = match side {
1071 OrderSide::Buy => add_price_ticks(order_price, increment, 1, precision),
1072 OrderSide::Sell => sub_price_ticks(order_price, increment, 1, precision),
1073 };
1074 let aggressive = clamp_price_to_range(aggressive, instrument, true);
1075 (aggressive != order_price).then_some(aggressive)
1076 }
1077
1078 fn limit_order_maintenance_suppressed(&self, side: OrderSide) -> bool {
1079 let state = self.limit_order_maintenance_state(side);
1080 matches!(
1081 state,
1082 LimitOrderMaintenanceState::AwaitingUpdate { .. }
1083 | LimitOrderMaintenanceState::AwaitingCancel { .. }
1084 | LimitOrderMaintenanceState::AwaitingReplacementAcceptance { .. }
1085 ) || (self.config.trigger_limit_order_maintenance_once
1086 && matches!(
1087 state,
1088 LimitOrderMaintenanceState::Completed | LimitOrderMaintenanceState::Failed
1089 ))
1090 }
1091
1092 fn limit_order_side_for_state(&self, client_order_id: ClientOrderId) -> Option<OrderSide> {
1093 [OrderSide::Buy, OrderSide::Sell].into_iter().find(|side| {
1094 self.limit_order_maintenance_state(*side).client_order_id() == Some(client_order_id)
1095 })
1096 }
1097
1098 const fn limit_order_maintenance_state(&self, side: OrderSide) -> LimitOrderMaintenanceState {
1099 match side {
1100 OrderSide::Buy => self.buy_limit_maintenance_state,
1101 OrderSide::Sell => self.sell_limit_maintenance_state,
1102 }
1103 }
1104
1105 fn set_limit_order_maintenance_state(
1106 &mut self,
1107 side: OrderSide,
1108 state: LimitOrderMaintenanceState,
1109 ) {
1110 match side {
1111 OrderSide::Buy => self.buy_limit_maintenance_state = state,
1112 OrderSide::Sell => self.sell_limit_maintenance_state = state,
1113 }
1114 }
1115
1116 fn tracked_limit_order(&self, side: OrderSide) -> Option<&OrderAny> {
1117 match side {
1118 OrderSide::Buy => self.buy_order.as_ref(),
1119 OrderSide::Sell => self.sell_order.as_ref(),
1120 }
1121 }
1122
1123 fn maintain_batch_limit_pair(&mut self, best_bid: Price, best_ask: Price) {
1125 self.refresh_tracked_order(OrderSide::Buy);
1129 self.refresh_tracked_order(OrderSide::Sell);
1130
1131 let Some(instrument) = &self.instrument else {
1132 return;
1133 };
1134 let Some(price_offset_ticks) = self.price_offset else {
1135 return;
1136 };
1137
1138 let buy_needs = match &self.buy_order {
1139 None => true,
1140 Some(order) => !Self::is_order_active(order) && !self.limit_order_is_one_shot(),
1141 };
1142 let sell_needs = match &self.sell_order {
1143 None => true,
1144 Some(order) => !Self::is_order_active(order) && !self.limit_order_is_one_shot(),
1145 };
1146
1147 if !buy_needs || !sell_needs {
1148 return;
1149 }
1150
1151 let increment = instrument.price_increment();
1152 let precision = instrument.price_precision();
1153
1154 let cross_spread = self.config.test_reject_post_only || self.config.limit_aggressive;
1158 let (unclamped_buy_price, unclamped_sell_price) = if cross_spread {
1159 (
1160 add_price_ticks(best_ask, increment, price_offset_ticks, precision),
1161 sub_price_ticks(best_bid, increment, price_offset_ticks, precision),
1162 )
1163 } else {
1164 (
1165 sub_price_ticks(best_bid, increment, price_offset_ticks, precision),
1166 add_price_ticks(best_ask, increment, price_offset_ticks, precision),
1167 )
1168 };
1169 let clamp = self.config.clamp_to_instrument_price_range;
1170 let buy_price = clamp_price_to_range(unclamped_buy_price, instrument, clamp);
1171 let sell_price = clamp_price_to_range(unclamped_sell_price, instrument, clamp);
1172 let quantity = instrument.make_qty(self.config.order_qty.as_f64(), None);
1173 let (time_in_force, expire_time) =
1174 self.resolve_time_in_force(self.config.limit_time_in_force);
1175 let instrument_id = self.config.instrument_id;
1176 let post_only = self.config.use_post_only || self.config.test_reject_post_only;
1177 let quote_quantity = self.config.use_quote_quantity;
1178 let display_qty = self.config.order_display_qty;
1179 let emulation_trigger = self.config.emulation_trigger;
1180
1181 let buy_order = self.order_factory().limit(
1182 instrument_id,
1183 OrderSide::Buy,
1184 quantity,
1185 buy_price,
1186 Some(time_in_force),
1187 expire_time,
1188 Some(post_only),
1189 None,
1190 Some(quote_quantity),
1191 display_qty,
1192 emulation_trigger,
1193 None,
1194 None,
1195 None,
1196 None,
1197 None,
1198 );
1199
1200 let sell_order = self.order_factory().limit(
1201 instrument_id,
1202 OrderSide::Sell,
1203 quantity,
1204 sell_price,
1205 Some(time_in_force),
1206 expire_time,
1207 Some(post_only),
1208 None,
1209 Some(quote_quantity),
1210 display_qty,
1211 emulation_trigger,
1212 None,
1213 None,
1214 None,
1215 None,
1216 None,
1217 );
1218
1219 self.buy_order = Some(buy_order.clone());
1220 self.sell_order = Some(sell_order.clone());
1221
1222 let client_id = self.config.client_id;
1223 if let Err(e) = self.submit_order_list(vec![buy_order, sell_order], None, client_id, None) {
1224 log::error!("Failed to submit batch limit pair: {e}");
1225 }
1226 }
1227
1228 fn maintain_stop_buy_orders(&mut self, best_bid: Price, best_ask: Price) {
1230 self.refresh_tracked_stop_order(OrderSide::Buy);
1231
1232 if let Some(order) = self.buy_stop_order.as_ref()
1234 && matches!(order.status(), OrderStatus::Rejected | OrderStatus::Denied)
1235 {
1236 return;
1237 }
1238
1239 let Some(instrument) = &self.instrument else {
1240 return;
1241 };
1242
1243 let increment = instrument.price_increment();
1244 let precision = instrument.price_precision();
1245 let stop_offset_ticks = self.config.stop_offset_ticks;
1246
1247 let unclamped_trigger_price = if matches!(
1249 self.config.stop_order_type,
1250 OrderType::LimitIfTouched | OrderType::MarketIfTouched | OrderType::TrailingStopMarket
1251 ) {
1252 sub_price_ticks(best_bid, increment, stop_offset_ticks, precision)
1254 } else {
1255 add_price_ticks(best_ask, increment, stop_offset_ticks, precision)
1257 };
1258 let clamp = self.config.clamp_to_instrument_price_range;
1259 let trigger_price = clamp_price_to_range(unclamped_trigger_price, instrument, clamp);
1260
1261 let limit_price = if matches!(
1263 self.config.stop_order_type,
1264 OrderType::StopLimit | OrderType::LimitIfTouched
1265 ) {
1266 let unclamped_limit_price =
1267 if let Some(limit_offset_ticks) = self.config.stop_limit_offset_ticks {
1268 add_price_ticks(trigger_price, increment, limit_offset_ticks, precision)
1270 } else {
1271 trigger_price
1272 };
1273 Some(clamp_price_to_range(
1274 unclamped_limit_price,
1275 instrument,
1276 clamp,
1277 ))
1278 } else {
1279 None
1280 };
1281
1282 let needs_new_order = match &self.buy_stop_order {
1283 None => true,
1284 Some(order) => !Self::is_order_active(order) && !self.stop_order_is_one_shot(),
1285 };
1286
1287 if needs_new_order {
1288 if let Err(e) = self.submit_stop_order(OrderSide::Buy, trigger_price, limit_price) {
1289 log::error!("Failed to submit buy stop order: {e}");
1290 }
1291 } else if let Some(order) = &self.buy_stop_order
1292 && order.venue_order_id().is_some()
1293 && !order.is_pending_update()
1294 && !order.is_pending_cancel()
1295 {
1296 let current_trigger = Self::get_order_trigger_price(order);
1297 if current_trigger.is_some() && current_trigger != Some(trigger_price) {
1298 if self.config.modify_stop_orders_to_maintain_offset {
1299 let order_clone = order.clone();
1300 if let Err(e) = self.modify_stop_order(&order_clone, trigger_price, limit_price)
1301 {
1302 log::error!("Failed to modify buy stop order: {e}");
1303 }
1304 } else if self.config.cancel_replace_stop_orders_to_maintain_offset
1305 && !self.buy_stop_cancel_replace_attempted
1306 {
1307 self.buy_stop_cancel_replace_attempted = true;
1308 let order_clone = order.clone();
1309 let _ = self.cancel_order(
1310 order_clone.client_order_id(),
1311 self.config.client_id,
1312 None,
1313 );
1314
1315 if let Err(e) =
1316 self.submit_stop_order(OrderSide::Buy, trigger_price, limit_price)
1317 {
1318 log::error!("Failed to submit replacement buy stop order: {e}");
1319 }
1320 }
1321 }
1322 }
1323 }
1324
1325 fn maintain_stop_sell_orders(&mut self, best_bid: Price, best_ask: Price) {
1327 self.refresh_tracked_stop_order(OrderSide::Sell);
1328
1329 if let Some(order) = self.sell_stop_order.as_ref()
1331 && matches!(order.status(), OrderStatus::Rejected | OrderStatus::Denied)
1332 {
1333 return;
1334 }
1335
1336 let Some(instrument) = &self.instrument else {
1337 return;
1338 };
1339
1340 let increment = instrument.price_increment();
1341 let precision = instrument.price_precision();
1342 let stop_offset_ticks = self.config.stop_offset_ticks;
1343
1344 let unclamped_trigger_price = if matches!(
1346 self.config.stop_order_type,
1347 OrderType::LimitIfTouched | OrderType::MarketIfTouched | OrderType::TrailingStopMarket
1348 ) {
1349 add_price_ticks(best_ask, increment, stop_offset_ticks, precision)
1351 } else {
1352 sub_price_ticks(best_bid, increment, stop_offset_ticks, precision)
1354 };
1355 let clamp = self.config.clamp_to_instrument_price_range;
1356 let trigger_price = clamp_price_to_range(unclamped_trigger_price, instrument, clamp);
1357
1358 let limit_price = if matches!(
1360 self.config.stop_order_type,
1361 OrderType::StopLimit | OrderType::LimitIfTouched
1362 ) {
1363 let unclamped_limit_price =
1364 if let Some(limit_offset_ticks) = self.config.stop_limit_offset_ticks {
1365 sub_price_ticks(trigger_price, increment, limit_offset_ticks, precision)
1367 } else {
1368 trigger_price
1369 };
1370 Some(clamp_price_to_range(
1371 unclamped_limit_price,
1372 instrument,
1373 clamp,
1374 ))
1375 } else {
1376 None
1377 };
1378
1379 let needs_new_order = match &self.sell_stop_order {
1380 None => true,
1381 Some(order) => !Self::is_order_active(order) && !self.stop_order_is_one_shot(),
1382 };
1383
1384 if needs_new_order {
1385 if let Err(e) = self.submit_stop_order(OrderSide::Sell, trigger_price, limit_price) {
1386 log::error!("Failed to submit sell stop order: {e}");
1387 }
1388 } else if let Some(order) = &self.sell_stop_order
1389 && order.venue_order_id().is_some()
1390 && !order.is_pending_update()
1391 && !order.is_pending_cancel()
1392 {
1393 let current_trigger = Self::get_order_trigger_price(order);
1394 if current_trigger.is_some() && current_trigger != Some(trigger_price) {
1395 if self.config.modify_stop_orders_to_maintain_offset {
1396 let order_clone = order.clone();
1397 if let Err(e) = self.modify_stop_order(&order_clone, trigger_price, limit_price)
1398 {
1399 log::error!("Failed to modify sell stop order: {e}");
1400 }
1401 } else if self.config.cancel_replace_stop_orders_to_maintain_offset
1402 && !self.sell_stop_cancel_replace_attempted
1403 {
1404 self.sell_stop_cancel_replace_attempted = true;
1405 let order_clone = order.clone();
1406 let _ = self.cancel_order(
1407 order_clone.client_order_id(),
1408 self.config.client_id,
1409 None,
1410 );
1411
1412 if let Err(e) =
1413 self.submit_stop_order(OrderSide::Sell, trigger_price, limit_price)
1414 {
1415 log::error!("Failed to submit replacement sell stop order: {e}");
1416 }
1417 }
1418 }
1419 }
1420 }
1421
1422 pub(super) fn submit_limit_order(
1428 &mut self,
1429 order_side: OrderSide,
1430 price: Price,
1431 ) -> anyhow::Result<()> {
1432 let Some(instrument) = self.instrument.as_ref() else {
1433 anyhow::bail!("No instrument loaded");
1434 };
1435 let quantity = instrument.make_qty(self.config.order_qty.as_f64(), None);
1436 let Some(order) = self.create_tracked_limit_order(order_side, price, quantity)? else {
1437 return Ok(());
1438 };
1439 let client_order_id = order.client_order_id();
1440 self.arm_limit_order_maintenance(order_side, client_order_id);
1441
1442 if let Err(e) = self.submit_order_apply_params(order) {
1443 self.handle_limit_order_closed(client_order_id);
1444 return Err(e);
1445 }
1446 Ok(())
1447 }
1448
1449 fn submit_limit_order_replacement(
1450 &mut self,
1451 order_side: OrderSide,
1452 price: Price,
1453 quantity: Quantity,
1454 ) -> anyhow::Result<()> {
1455 let Some(order) = self.create_tracked_limit_order(order_side, price, quantity)? else {
1456 anyhow::bail!("Replacement {order_side} limit order was not created");
1457 };
1458 let client_order_id = order.client_order_id();
1459 self.set_limit_order_maintenance_state(
1460 order_side,
1461 LimitOrderMaintenanceState::AwaitingReplacementAcceptance {
1462 client_order_id,
1463 price,
1464 },
1465 );
1466 self.submit_order_apply_params(order)
1467 }
1468
1469 fn create_tracked_limit_order(
1470 &mut self,
1471 order_side: OrderSide,
1472 price: Price,
1473 quantity: Quantity,
1474 ) -> anyhow::Result<Option<OrderAny>> {
1475 if self.instrument.is_none() {
1476 anyhow::bail!("No instrument loaded");
1477 }
1478
1479 if self.config.dry_run {
1480 log_warn!("Dry run, skipping create {order_side:?} order");
1481 return Ok(None);
1482 }
1483
1484 if order_side == OrderSide::Buy && !self.config.enable_limit_buys {
1485 log_warn!("BUY orders not enabled, skipping");
1486 return Ok(None);
1487 } else if order_side == OrderSide::Sell && !self.config.enable_limit_sells {
1488 log_warn!("SELL orders not enabled, skipping");
1489 return Ok(None);
1490 }
1491
1492 let (time_in_force, expire_time) =
1493 self.resolve_time_in_force(self.config.limit_time_in_force);
1494
1495 let instrument_id = self.config.instrument_id;
1496 let post_only = self.config.use_post_only || self.config.test_reject_post_only;
1497 let quote_quantity = self.config.use_quote_quantity;
1498 let display_qty = self
1499 .config
1500 .order_display_qty
1501 .map(|display_qty| display_qty.min(quantity));
1502 let emulation_trigger = self.config.emulation_trigger;
1503
1504 let order = self.order_factory().limit(
1505 instrument_id,
1506 order_side,
1507 quantity,
1508 price,
1509 Some(time_in_force),
1510 expire_time,
1511 Some(post_only),
1512 None, Some(quote_quantity),
1514 display_qty,
1515 emulation_trigger,
1516 None, None, None, None, None, );
1522
1523 if order_side == OrderSide::Buy {
1524 self.buy_order = Some(order.clone());
1525 } else {
1526 self.sell_order = Some(order.clone());
1527 }
1528
1529 Ok(Some(order))
1530 }
1531
1532 fn arm_limit_order_maintenance(&mut self, side: OrderSide, client_order_id: ClientOrderId) {
1533 if matches!(
1534 self.limit_order_maintenance_state(side),
1535 LimitOrderMaintenanceState::Ready
1536 | LimitOrderMaintenanceState::AwaitingAcceptance { .. }
1537 | LimitOrderMaintenanceState::Active { .. }
1538 ) {
1539 self.set_limit_order_maintenance_state(
1540 side,
1541 LimitOrderMaintenanceState::AwaitingAcceptance { client_order_id },
1542 );
1543 }
1544 }
1545
1546 #[expect(
1552 clippy::too_many_lines,
1553 reason = "stop order submission covers all supported stop order scenarios"
1554 )]
1555 pub(super) fn submit_stop_order(
1556 &mut self,
1557 order_side: OrderSide,
1558 trigger_price: Price,
1559 limit_price: Option<Price>,
1560 ) -> anyhow::Result<()> {
1561 let Some(instrument) = &self.instrument else {
1562 anyhow::bail!("No instrument loaded");
1563 };
1564
1565 if self.config.dry_run {
1566 log_warn!("Dry run, skipping create {order_side:?} stop order");
1567 return Ok(());
1568 }
1569
1570 if order_side == OrderSide::Buy && !self.config.enable_stop_buys {
1571 log_warn!("BUY stop orders not enabled, skipping");
1572 return Ok(());
1573 } else if order_side == OrderSide::Sell && !self.config.enable_stop_sells {
1574 log_warn!("SELL stop orders not enabled, skipping");
1575 return Ok(());
1576 }
1577
1578 let (time_in_force, expire_time) =
1579 self.resolve_time_in_force(self.config.stop_time_in_force);
1580
1581 let quantity = instrument.make_qty(self.config.order_qty.as_f64(), None);
1583 let instrument_id = self.config.instrument_id;
1584 let trigger_type = self.config.stop_trigger_type;
1585 let quote_quantity = self.config.use_quote_quantity;
1586 let display_qty = self.config.order_display_qty;
1587 let emulation_trigger = self.config.emulation_trigger;
1588 let stop_order_type = self.config.stop_order_type;
1589 let trailing_offset = self.config.trailing_offset;
1590 let trailing_offset_type = self.config.trailing_offset_type;
1591
1592 let mut factory = self.order_factory();
1593
1594 let mut order: OrderAny = match stop_order_type {
1595 OrderType::StopMarket => factory.stop_market(
1596 instrument_id,
1597 order_side,
1598 quantity,
1599 trigger_price,
1600 Some(trigger_type),
1601 Some(time_in_force),
1602 expire_time,
1603 None, Some(quote_quantity),
1605 None, emulation_trigger,
1607 None, None, None, None, None, ),
1613 OrderType::StopLimit => {
1614 let Some(limit_price) = limit_price else {
1615 anyhow::bail!("STOP_LIMIT order requires limit_price");
1616 };
1617 factory.stop_limit(
1618 instrument_id,
1619 order_side,
1620 quantity,
1621 limit_price,
1622 trigger_price,
1623 Some(trigger_type),
1624 Some(time_in_force),
1625 expire_time,
1626 None, None, Some(quote_quantity),
1629 display_qty,
1630 emulation_trigger,
1631 None, None, None, None, None, )
1637 }
1638 OrderType::MarketIfTouched => factory.market_if_touched(
1639 instrument_id,
1640 order_side,
1641 quantity,
1642 trigger_price,
1643 Some(trigger_type),
1644 Some(time_in_force),
1645 expire_time,
1646 None, Some(quote_quantity),
1648 emulation_trigger,
1649 None, None, None, None, None, ),
1655 OrderType::LimitIfTouched => {
1656 let Some(limit_price) = limit_price else {
1657 anyhow::bail!("LIMIT_IF_TOUCHED order requires limit_price");
1658 };
1659 factory.limit_if_touched(
1660 instrument_id,
1661 order_side,
1662 quantity,
1663 limit_price,
1664 trigger_price,
1665 Some(trigger_type),
1666 Some(time_in_force),
1667 expire_time,
1668 None, None, Some(quote_quantity),
1671 display_qty,
1672 emulation_trigger,
1673 None, None, None, None, None, )
1679 }
1680 OrderType::TrailingStopMarket => {
1681 let Some(trailing_offset) = trailing_offset else {
1682 anyhow::bail!("TRAILING_STOP_MARKET order requires trailing_offset config");
1683 };
1684 factory.trailing_stop_market(
1685 instrument_id,
1686 order_side,
1687 quantity,
1688 trailing_offset,
1689 Some(trailing_offset_type),
1690 None,
1691 Some(trigger_price),
1692 Some(trigger_type),
1693 Some(time_in_force),
1694 expire_time,
1695 None, Some(quote_quantity),
1697 None, emulation_trigger,
1699 None, None, None, None, None, )
1705 }
1706 _ => {
1707 anyhow::bail!("Unknown stop order type: {stop_order_type:?}");
1708 }
1709 };
1710 drop(factory);
1711
1712 if let OrderAny::TrailingStopMarket(order) = &mut order {
1713 order.activation_price = Some(trigger_price);
1714 }
1715
1716 if order_side == OrderSide::Buy {
1717 self.buy_stop_order = Some(order.clone());
1718 } else {
1719 self.sell_stop_order = Some(order.clone());
1720 }
1721
1722 self.submit_order_apply_params(order)
1723 }
1724
1725 pub(super) fn submit_bracket_order(
1731 &mut self,
1732 order_side: OrderSide,
1733 entry_price: Price,
1734 ) -> anyhow::Result<()> {
1735 let Some(instrument) = &self.instrument else {
1736 anyhow::bail!("No instrument loaded");
1737 };
1738
1739 if self.config.dry_run {
1740 log_warn!("Dry run, skipping create {order_side:?} bracket order");
1741 return Ok(());
1742 }
1743
1744 if self.config.bracket_entry_order_type != OrderType::Limit {
1745 anyhow::bail!(
1746 "Only Limit entry orders are supported for brackets, was {:?}",
1747 self.config.bracket_entry_order_type
1748 );
1749 }
1750
1751 if order_side == OrderSide::Buy && !self.config.enable_limit_buys {
1752 log_warn!("BUY orders not enabled, skipping bracket");
1753 return Ok(());
1754 } else if order_side == OrderSide::Sell && !self.config.enable_limit_sells {
1755 log_warn!("SELL orders not enabled, skipping bracket");
1756 return Ok(());
1757 }
1758
1759 let (time_in_force, expire_time) =
1760 self.resolve_time_in_force(self.config.limit_time_in_force);
1761 let sl_time_in_force = self.config.stop_time_in_force.unwrap_or(TimeInForce::Gtc);
1762 if sl_time_in_force == TimeInForce::Gtd {
1763 anyhow::bail!("GTD time in force not supported for bracket stop-loss legs");
1764 }
1765
1766 let quantity = instrument.make_qty(self.config.order_qty.as_f64(), None);
1767 let increment = instrument.price_increment();
1768 let precision = instrument.price_precision();
1769 let bracket_offset_ticks = self.config.bracket_offset_ticks;
1770
1771 let (unclamped_tp_price, unclamped_sl_trigger_price) = match order_side {
1772 OrderSide::Buy => {
1773 let tp = add_price_ticks(entry_price, increment, bracket_offset_ticks, precision);
1774 let sl = sub_price_ticks(entry_price, increment, bracket_offset_ticks, precision);
1775 (tp, sl)
1776 }
1777 OrderSide::Sell => {
1778 let tp = sub_price_ticks(entry_price, increment, bracket_offset_ticks, precision);
1779 let sl = add_price_ticks(entry_price, increment, bracket_offset_ticks, precision);
1780 (tp, sl)
1781 }
1782 };
1783 let clamp = self.config.clamp_to_instrument_price_range;
1784 let tp_price = clamp_price_to_range(unclamped_tp_price, instrument, clamp);
1785 let sl_trigger_price = clamp_price_to_range(unclamped_sl_trigger_price, instrument, clamp);
1786
1787 let entry_post_only = self.config.use_post_only || self.config.test_reject_post_only;
1788 let instrument_id = self.config.instrument_id;
1789 let quote_quantity = self.config.use_quote_quantity;
1790 let emulation_trigger = self.config.emulation_trigger;
1791 let stop_trigger_type = self.config.stop_trigger_type;
1792 let orders = self
1793 .order_factory()
1794 .bracket()
1795 .instrument_id(instrument_id)
1796 .order_side(order_side)
1797 .quantity(quantity)
1798 .quote_quantity(quote_quantity)
1799 .entry_order_type(OrderType::Limit)
1800 .entry_price(entry_price)
1801 .time_in_force(time_in_force)
1802 .entry_post_only(entry_post_only)
1803 .maybe_emulation_trigger(emulation_trigger)
1804 .maybe_expire_time(expire_time)
1805 .tp_price(tp_price)
1806 .tp_post_only(entry_post_only)
1807 .tp_time_in_force(time_in_force)
1808 .sl_trigger_price(sl_trigger_price)
1809 .sl_trigger_type(stop_trigger_type)
1810 .sl_time_in_force(sl_time_in_force)
1811 .call();
1812
1813 if let Some(entry_order) = orders.first() {
1814 if order_side == OrderSide::Buy {
1815 self.buy_order = Some(entry_order.clone());
1816 } else {
1817 self.sell_order = Some(entry_order.clone());
1818 }
1819 self.arm_limit_order_maintenance(order_side, entry_order.client_order_id());
1820 }
1821
1822 let client_id = self.config.client_id;
1823 if let Some(params) = &self.config.order_params {
1824 self.submit_order_list(orders, None, client_id, Some(params.clone()))
1825 } else {
1826 self.submit_order_list(orders, None, client_id, None)
1827 }
1828 }
1829
1830 fn close_positions_on_stop(
1831 &mut self,
1832 strategy_id: StrategyId,
1833 time_in_force: Option<TimeInForce>,
1834 ) -> anyhow::Result<()> {
1835 let instrument_id = self.config.instrument_id;
1836 let client_id = self.config.client_id;
1837 let reduce_only = Some(self.config.reduce_only_on_stop);
1838 let Some(precision) = self.config.close_positions_qty_precision else {
1839 return self.close_all_positions(
1840 instrument_id,
1841 None,
1842 client_id,
1843 None,
1844 time_in_force,
1845 reduce_only,
1846 None,
1847 None,
1848 );
1849 };
1850
1851 let positions =
1852 self.cache()
1853 .positions_open(None, Some(&instrument_id), Some(&strategy_id), None, None);
1854
1855 if positions.is_empty() {
1856 log::info!("No {instrument_id} open positions to close");
1857 return Ok(());
1858 }
1859
1860 log::info!(
1861 "Closing {} open position{}",
1862 positions.len(),
1863 if positions.len() == 1 { "" } else { "s" },
1864 );
1865
1866 for position in positions {
1867 let position_id = position.id;
1868 let (close_quantity, residual) = split_position_quantity(position.quantity, precision)?;
1869
1870 if close_quantity.is_zero() {
1871 log_warn!(
1872 "Position {position_id} has no venue-fillable close quantity at {precision}-decimal precision; exact residual remains {residual}"
1873 );
1874 continue;
1875 }
1876
1877 if residual > Decimal::ZERO {
1878 log_warn!(
1879 "Position {position_id} close quantity {close_quantity} leaves exact residual {residual} after a full fill"
1880 );
1881 }
1882
1883 let Some(closing_side) = OrderCore::closing_side(position.side) else {
1884 continue;
1885 };
1886 let order = self.order_factory().market(
1887 position.instrument_id,
1888 closing_side,
1889 close_quantity,
1890 time_in_force,
1891 reduce_only,
1892 None,
1893 None,
1894 None,
1895 None,
1896 None,
1897 );
1898
1899 self.submit_order(order, Some(position_id), client_id, None)?;
1900 }
1901
1902 Ok(())
1903 }
1904
1905 pub(super) fn open_position(&mut self, net_qty: Decimal) -> anyhow::Result<()> {
1911 let Some(instrument) = &self.instrument else {
1912 anyhow::bail!("No instrument loaded");
1913 };
1914
1915 if self.config.dry_run {
1916 log_warn!("Dry run, skipping open position");
1917 return Ok(());
1918 }
1919
1920 if net_qty == Decimal::ZERO {
1921 log_warn!("Open position with zero quantity, skipping");
1922 return Ok(());
1923 }
1924
1925 let order_side = if net_qty > Decimal::ZERO {
1926 OrderSide::Buy
1927 } else {
1928 OrderSide::Sell
1929 };
1930
1931 let quantity = instrument.make_qty(net_qty.abs().to_f64().unwrap_or(0.0), None);
1932
1933 let reduce_only = if self.config.test_reject_reduce_only {
1935 Some(true)
1936 } else {
1937 None
1938 };
1939 let instrument_id = self.config.instrument_id;
1940 let time_in_force = self.config.open_position_time_in_force;
1941 let quote_quantity = self.config.use_quote_quantity;
1942
1943 let order = self.order_factory().market(
1944 instrument_id,
1945 order_side,
1946 quantity,
1947 Some(time_in_force),
1948 reduce_only,
1949 Some(quote_quantity),
1950 None, None, None, None, );
1955
1956 self.submit_order_apply_params(order)
1957 }
1958
1959 pub(super) fn cancel_active_orders(
1960 &mut self,
1961 instrument_id: InstrumentId,
1962 strategy_id: StrategyId,
1963 client_id: Option<ClientId>,
1964 ) {
1965 let bracket_targets: Vec<ClientOrderId> = {
1968 let cache = self.cache();
1969 let mut targets = Vec::new();
1970
1971 for order_list in
1972 cache.order_lists(None, Some(&instrument_id), Some(&strategy_id), None)
1973 {
1974 let is_bracket = order_list.client_order_ids.iter().any(|cid| {
1975 cache
1976 .order(cid)
1977 .is_some_and(|o| is_in_contingency_group(&o))
1978 });
1979
1980 if !is_bracket {
1981 continue;
1982 }
1983
1984 for cid in &order_list.client_order_ids {
1985 if let Some(order) = cache.order(cid)
1986 && !order.is_closed()
1987 && !order.is_pending_cancel()
1988 {
1989 targets.push(*cid);
1990 }
1991 }
1992 }
1993 targets
1994 };
1995
1996 for cid in bracket_targets {
1997 if let Err(e) = self.cancel_order(cid, client_id, None) {
1998 log::error!("Failed to cancel bracket leg {cid}: {e}");
1999 }
2000 }
2001
2002 if self.config.use_individual_cancels_on_stop {
2003 for cid in self.collect_cancellable_order_ids(instrument_id, strategy_id) {
2004 if let Err(e) = self.cancel_order(cid, client_id, None) {
2005 log::error!("Failed to cancel order {cid}: {e}");
2006 }
2007 }
2008 } else if self.config.use_batch_cancel_on_stop {
2009 let candidates = self.collect_cancellable_orders(instrument_id, strategy_id);
2010 let mut batchable: Vec<ClientOrderId> = Vec::new();
2011
2012 for order in candidates {
2013 let cid = order.client_order_id();
2014 if order.is_emulated() || order.is_active_local() {
2015 if let Err(e) = self.cancel_order(cid, client_id, None) {
2016 log::error!("Failed to cancel local order {cid}: {e}");
2017 }
2018 } else {
2019 batchable.push(cid);
2020 }
2021 }
2022
2023 if !batchable.is_empty()
2024 && let Err(e) = self.cancel_orders(batchable, client_id, None)
2025 {
2026 log::error!("Failed to batch cancel orders: {e}");
2027 }
2028 } else {
2029 let local_ids: Vec<ClientOrderId> = {
2032 let cache = self.cache();
2033 cache
2034 .orders_active_local(None, Some(&instrument_id), Some(&strategy_id), None, None)
2035 .into_iter()
2036 .filter(|o| {
2037 !o.is_closed() && !o.is_pending_cancel() && !is_in_contingency_group(o)
2038 })
2039 .map(|o| o.client_order_id())
2040 .collect()
2041 };
2042
2043 for cid in local_ids {
2044 if let Err(e) = self.cancel_order(cid, client_id, None) {
2045 log::error!("Failed to cancel active-local order {cid}: {e}");
2046 }
2047 }
2048
2049 if let Err(e) = self.cancel_all_orders(instrument_id, None, client_id, false, None) {
2050 log::error!("Failed to cancel all orders: {e}");
2051 }
2052 }
2053 }
2054
2055 pub(super) fn collect_cancellable_orders(
2056 &self,
2057 instrument_id: InstrumentId,
2058 strategy_id: StrategyId,
2059 ) -> Vec<OrderAny> {
2060 let cache = self.cache();
2061 let mut seen: AHashSet<ClientOrderId> = AHashSet::new();
2062 let mut candidates: Vec<OrderAny> = Vec::new();
2063 let sources = [
2066 cache.orders_active_local(None, Some(&instrument_id), Some(&strategy_id), None, None),
2067 cache.orders_emulated(None, Some(&instrument_id), Some(&strategy_id), None, None),
2068 cache.orders_inflight(None, Some(&instrument_id), Some(&strategy_id), None, None),
2069 cache.orders_open(None, Some(&instrument_id), Some(&strategy_id), None, None),
2070 ];
2071
2072 for orders in sources {
2073 for order in orders {
2074 if order.is_closed() || order.is_pending_cancel() || is_in_contingency_group(&order)
2075 {
2076 continue;
2077 }
2078 let cid = order.client_order_id();
2079 if seen.insert(cid) {
2080 candidates.push(order);
2081 }
2082 }
2083 }
2084 candidates
2085 }
2086
2087 pub(super) fn collect_cancellable_order_ids(
2088 &self,
2089 instrument_id: InstrumentId,
2090 strategy_id: StrategyId,
2091 ) -> Vec<ClientOrderId> {
2092 self.collect_cancellable_orders(instrument_id, strategy_id)
2093 .into_iter()
2094 .map(|o| o.client_order_id())
2095 .collect()
2096 }
2097}
2098
2099impl LimitOrderMaintenanceState {
2100 const fn new(enabled: bool) -> Self {
2101 if enabled { Self::Ready } else { Self::Disabled }
2102 }
2103
2104 const fn client_order_id(self) -> Option<ClientOrderId> {
2105 match self {
2106 Self::AwaitingAcceptance { client_order_id }
2107 | Self::Active { client_order_id }
2108 | Self::AwaitingUpdate {
2109 client_order_id, ..
2110 }
2111 | Self::AwaitingCancel {
2112 client_order_id, ..
2113 }
2114 | Self::AwaitingReplacementAcceptance {
2115 client_order_id, ..
2116 } => Some(client_order_id),
2117 Self::Disabled | Self::Ready | Self::Completed | Self::Failed => None,
2118 }
2119 }
2120}
2121
2122fn split_position_quantity(
2123 quantity: Quantity,
2124 precision: u8,
2125) -> anyhow::Result<(Quantity, Decimal)> {
2126 let quantity_decimal = quantity.as_decimal();
2127 let close_decimal = quantity_decimal.trunc_with_scale(u32::from(precision));
2128 let close_quantity = Quantity::from_decimal_dp(close_decimal, quantity.precision)
2129 .map_err(|e| anyhow::anyhow!("Invalid position close quantity {close_decimal}: {e}"))?;
2130 let residual = quantity_decimal - close_quantity.as_decimal();
2131 Ok((close_quantity, residual))
2132}
2133
2134fn add_price_ticks(base: Price, increment: Price, ticks: u64, precision: u8) -> Price {
2135 let offset = price_tick_offset(increment, ticks, precision);
2136 base + offset
2137}
2138
2139fn sub_price_ticks(base: Price, increment: Price, ticks: u64, precision: u8) -> Price {
2140 let offset = price_tick_offset(increment, ticks, precision);
2141 base - offset
2142}
2143
2144fn price_tick_offset(increment: Price, ticks: u64, precision: u8) -> Price {
2145 let offset = increment * Decimal::from(ticks);
2146 Price::from_decimal_dp(offset, precision)
2147 .unwrap_or_else(|e| panic!("Failed to calculate price tick offset: {e}"))
2148}
2149
2150fn is_in_contingency_group(order: &OrderAny) -> bool {
2151 matches!(
2152 order.contingency_type(),
2153 Some(ContingencyType::Oto | ContingencyType::Oco | ContingencyType::Ouo)
2154 )
2155}
2156
2157fn clamp_price_to_range(price: Price, instrument: &InstrumentAny, enabled: bool) -> Price {
2158 if !enabled {
2159 return price;
2160 }
2161 let mut clamped = price;
2162 if let Some(min) = instrument.min_price()
2163 && clamped < min
2164 {
2165 clamped = min;
2166 }
2167
2168 if let Some(max) = instrument.max_price()
2169 && clamped > max
2170 {
2171 clamped = max;
2172 }
2173
2174 clamped
2175}