1#[path = "core_orders.rs"]
19mod core_orders;
20#[path = "core_tracking.rs"]
21mod core_tracking;
22#[path = "core_updates.rs"]
23mod core_updates;
24#[cfg(test)]
25#[path = "core_tests.rs"]
26mod tests;
27
28use std::{
29 collections::VecDeque,
30 fmt::Debug,
31 str::FromStr,
32 sync::{
33 Arc,
34 atomic::{AtomicBool, Ordering},
35 },
36 time::Duration,
37};
38
39use ahash::AHashMap;
40use anyhow::Context;
41use ibapi::{
42 accounts::PositionUpdate,
43 client::Client,
44 contracts::{Contract, SecurityType},
45 orders::{
46 ExecutionData, ExecutionFilter, Executions, OrderStatus as IBOrderStatus, OrderUpdate,
47 Orders,
48 },
49 prelude::{StreamExt, SubscriptionItemStreamExt},
50};
51use nautilus_common::{
52 cache::{Cache, fifo::FifoCacheMap},
53 clients::ExecutionClient,
54 enums::LogLevel,
55 factories::OrderEventFactory,
56 live::{runner::get_exec_event_sender, sender::EventSender},
57 messages::{
58 ExecutionEvent,
59 execution::{
60 BatchCancelOrders, CancelAllOrders, CancelOrder, ExecutionReport, GenerateFillReports,
61 GenerateFillReportsBuilder, GenerateOrderStatusReport, GenerateOrderStatusReports,
62 GenerateOrderStatusReportsBuilder, GeneratePositionStatusReports,
63 GeneratePositionStatusReportsBuilder, ModifyOrder, QueryAccount, QueryOrder,
64 SubmitOrder, SubmitOrderList,
65 },
66 },
67 msgbus::{send_account_state, switchboard::MessagingSwitchboard},
68};
69use nautilus_core::{
70 DurationNanos, Params, UUID4, UnixNanos,
71 time::{AtomicTime, get_atomic_clock_realtime},
72};
73use nautilus_live::{
74 ExecutionClientCore,
75 execution::failure::CommandFailure,
76 task::{TaskGroup, TaskGroupGuard},
77};
78use nautilus_model::{
79 accounts::AccountAny,
80 enums::{
81 LiquiditySide, OmsType, OrderSide, OrderStatus, OrderType, PositionSide, TimeInForce,
82 TrailingOffsetType,
83 },
84 events::{
85 AccountState, OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied,
86 OrderDeniedReason, OrderEventAny, OrderFilled, OrderModifyRejected, OrderPendingCancel,
87 OrderRejected, OrderSubmitted, OrderUpdated,
88 },
89 identifiers::{
90 AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, TradeId, TraderId, Venue,
91 VenueOrderId,
92 },
93 instruments::{Instrument, InstrumentAny},
94 orders::{Order, any::OrderAny},
95 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
96 types::{AccountBalance, Currency, MarginBalance, Money, Price, Quantity},
97};
98use parking_lot::Mutex;
99use rust_decimal::{Decimal, prelude::ToPrimitive};
100use ustr::Ustr;
101
102use super::{
103 account::{PositionTracker, create_position_tracker, raw_ib_account_code},
104 parse::{
105 ib_venue_order_id, parse_execution_time, parse_execution_to_fill_report,
106 parse_order_status_to_report,
107 },
108 transform::nautilus_order_to_ib_order,
109};
110use crate::{
111 common::{
112 parse::{ib_contract_to_instrument_id_simple, is_spread_instrument_id},
113 shared_client::SharedClientHandle,
114 },
115 config::InteractiveBrokersExecutionClientConfig,
116 providers::instruments::InteractiveBrokersInstrumentProvider,
117};
118
119#[cfg_attr(
124 feature = "python",
125 pyo3::pyclass(module = "nautilus_trader.adapters.interactive_brokers", unsendable)
126)]
127pub struct InteractiveBrokersExecutionClient {
128 core: ExecutionClientCore,
129 config: InteractiveBrokersExecutionClientConfig,
130 instrument_provider: Arc<InteractiveBrokersInstrumentProvider>,
131 is_connected: AtomicBool,
132 ib_client: Option<SharedClientHandle>,
133 pending_tasks: TaskGroup,
134 next_order_id: Arc<Mutex<i32>>,
135 order_submit_lock: Arc<tokio::sync::Mutex<()>>,
136 session_tasks: TaskGroup,
137 order_id_map: Arc<Mutex<AHashMap<ClientOrderId, i32>>>,
138 venue_order_id_map: Arc<Mutex<AHashMap<i32, ClientOrderId>>>,
139 commission_cache: Arc<Mutex<CommissionCache>>,
140 pending_execution_cache: Arc<Mutex<PendingExecutionCache>>,
141 instrument_id_map: Arc<Mutex<AHashMap<i32, InstrumentId>>>,
142 trader_id_map: Arc<Mutex<AHashMap<i32, TraderId>>>,
143 strategy_id_map: Arc<Mutex<AHashMap<i32, StrategyId>>>,
144 active_order_contexts: Arc<Mutex<AHashMap<i32, TrackedOrderContext>>>,
145 terminal_order_contexts: Arc<Mutex<FifoCacheMap<i32, TrackedOrderContext, 10_000>>>,
146 spread_fill_tracking: Arc<Mutex<AHashMap<ClientOrderId, ahash::AHashSet<String>>>>,
147 position_tracker: PositionTracker,
148 order_avg_prices: Arc<Mutex<AHashMap<ClientOrderId, Price>>>,
149 pending_combo_fills: Arc<Mutex<AHashMap<ClientOrderId, VecDeque<PendingComboFill>>>>,
150 pending_combo_fill_avgs: Arc<Mutex<AHashMap<ClientOrderId, VecDeque<(Decimal, Price)>>>>,
151 order_fill_progress: Arc<Mutex<AHashMap<ClientOrderId, (Decimal, Decimal)>>>,
152 pending_cancel_orders: Arc<Mutex<ahash::AHashSet<ClientOrderId>>>,
153}
154
155type CommissionCache = FifoCacheMap<String, (f64, String), 10_000>;
156type PendingExecutionCache = FifoCacheMap<String, ExecutionData, 10_000>;
157
158#[derive(Clone, Debug)]
159struct PendingComboFill {
160 trader_id: TraderId,
161 strategy_id: StrategyId,
162 account_id: AccountId,
163 instrument_id: InstrumentId,
164 venue_order_id: VenueOrderId,
165 trade_id: TradeId,
166 order_side: OrderSide,
167 order_type: OrderType,
168 last_qty: Quantity,
169 commission: Money,
170 liquidity_side: LiquiditySide,
171 quote_currency: Currency,
172 client_order_id: ClientOrderId,
173 ts_event: UnixNanos,
174 ts_init: UnixNanos,
175}
176
177#[derive(Clone, Debug)]
178struct TrackedOrderContext {
179 client_order_id: ClientOrderId,
180 trader_id: TraderId,
181 strategy_id: StrategyId,
182 instrument_id: InstrumentId,
183 order_side: OrderSide,
184 order_type: OrderType,
185 accepted: bool,
186 avg_px: Option<Price>,
187}
188
189#[derive(Debug, Clone, Copy, PartialEq, Eq)]
190enum IbOrderSelector {
191 OrderId(i32),
192 PermId(i64),
193}
194
195impl IbOrderSelector {
196 fn from_venue_order_id(venue_order_id: &VenueOrderId) -> anyhow::Result<Self> {
197 let raw = venue_order_id.as_str();
198 if let Some(perm_id) = raw.strip_prefix("PERM-") {
199 return Ok(Self::PermId(perm_id.parse::<i64>().with_context(|| {
200 format!("Failed to parse venue_order_id {raw:?} as IB perm_id")
201 })?));
202 }
203
204 Ok(Self::OrderId(raw.parse::<i32>().with_context(|| {
205 format!("Failed to parse venue_order_id {raw:?} as IB order_id")
206 })?))
207 }
208
209 fn matches(self, order_id: i32, perm_id: i64) -> bool {
210 match self {
211 Self::OrderId(target_order_id) => order_id == target_order_id,
212 Self::PermId(target_perm_id) => perm_id == target_perm_id,
213 }
214 }
215
216 fn venue_order_id(self) -> VenueOrderId {
217 match self {
218 Self::OrderId(order_id) => VenueOrderId::from(order_id.to_string()),
219 Self::PermId(perm_id) => VenueOrderId::from(format!("PERM-{perm_id}")),
220 }
221 }
222
223 fn label(self) -> String {
224 self.venue_order_id().to_string()
225 }
226}
227
228impl Debug for InteractiveBrokersExecutionClient {
229 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
230 f.debug_struct(stringify!(InteractiveBrokersExecutionClient))
231 .field("core", &self.core)
232 .field("config", &self.config)
233 .field("instrument_provider", &self.instrument_provider)
234 .field("is_connected", &self.is_connected.load(Ordering::Relaxed))
235 .field("ib_client", &self.ib_client.is_some())
236 .finish_non_exhaustive()
237 }
238}
239
240impl InteractiveBrokersExecutionClient {
241 pub fn new(
253 mut core: ExecutionClientCore,
254 config: InteractiveBrokersExecutionClientConfig,
255 instrument_provider: Arc<InteractiveBrokersInstrumentProvider>,
256 ) -> anyhow::Result<Self> {
257 anyhow::ensure!(
258 !config.client_id.unsigned_abs().is_multiple_of(1000),
259 "Interactive Brokers execution client_id must not be a multiple of 1000 because order ID partitioning uses client_id % 1000; got {}",
260 config.client_id
261 );
262
263 if let Some(account_id) = &config.account_id {
265 core.account_id = AccountId::from(account_id.clone());
266 }
267
268 let pending_tasks = TaskGroup::new();
269 let session_tasks = TaskGroup::new();
270
271 Ok(Self {
272 core,
273 config,
274 instrument_provider,
275 is_connected: AtomicBool::new(false),
276 ib_client: None,
277 pending_tasks,
278 next_order_id: Arc::new(Mutex::new(0)),
279 order_submit_lock: Arc::new(tokio::sync::Mutex::new(())),
280 session_tasks,
281 order_id_map: Arc::new(Mutex::new(AHashMap::new())),
282 venue_order_id_map: Arc::new(Mutex::new(AHashMap::new())),
283 commission_cache: Arc::new(Mutex::new(CommissionCache::new())),
284 pending_execution_cache: Arc::new(Mutex::new(PendingExecutionCache::new())),
285 instrument_id_map: Arc::new(Mutex::new(AHashMap::new())),
286 trader_id_map: Arc::new(Mutex::new(AHashMap::new())),
287 strategy_id_map: Arc::new(Mutex::new(AHashMap::new())),
288 active_order_contexts: Arc::new(Mutex::new(AHashMap::new())),
289 terminal_order_contexts: Arc::new(Mutex::new(FifoCacheMap::new())),
290 spread_fill_tracking: Arc::new(Mutex::new(AHashMap::new())),
291 position_tracker: create_position_tracker(),
292 order_avg_prices: Arc::new(Mutex::new(AHashMap::new())),
293 pending_combo_fills: Arc::new(Mutex::new(AHashMap::new())),
294 pending_combo_fill_avgs: Arc::new(Mutex::new(AHashMap::new())),
295 order_fill_progress: Arc::new(Mutex::new(AHashMap::new())),
296 pending_cancel_orders: Arc::new(Mutex::new(ahash::AHashSet::new())),
297 })
298 }
299
300 fn submit_order_list_with_orders(
301 &self,
302 cmd: SubmitOrderList,
303 orders: Vec<OrderAny>,
304 ) -> anyhow::Result<()> {
305 let client = self.ib_client.as_ref().context("IB client not connected")?;
306
307 let order_id_map = Arc::clone(&self.order_id_map);
308 let venue_order_id_map = Arc::clone(&self.venue_order_id_map);
309 let instrument_id_map = Arc::clone(&self.instrument_id_map);
310 let trader_id_map = Arc::clone(&self.trader_id_map);
311 let strategy_id_map = Arc::clone(&self.strategy_id_map);
312 let active_order_contexts = Arc::clone(&self.active_order_contexts);
313 let terminal_order_contexts = Arc::clone(&self.terminal_order_contexts);
314 let next_order_id = Arc::clone(&self.next_order_id);
315 let instrument_provider = Arc::clone(&self.instrument_provider);
316 let exec_sender = get_exec_event_sender();
317 let clock = get_atomic_clock_realtime();
318 let account_id = self.core.account_id;
319 let strategy_id = cmd.strategy_id;
320 let client_clone = client.as_arc().clone();
321 let order_submit_lock = Arc::clone(&self.order_submit_lock);
322
323 let future = async move {
324 if let Err(e) = Self::handle_submit_order_list_async(
325 &cmd,
326 &orders,
327 &client_clone,
328 &order_id_map,
329 &venue_order_id_map,
330 &instrument_id_map,
331 &trader_id_map,
332 &strategy_id_map,
333 &active_order_contexts,
334 &terminal_order_contexts,
335 &next_order_id,
336 &instrument_provider,
337 &exec_sender,
338 clock,
339 account_id,
340 strategy_id,
341 &order_submit_lock,
342 )
343 .await
344 {
345 tracing::error!("Error submitting order list: {e}");
346 }
347 };
348
349 self.pending_tasks
350 .spawn(future)
351 .context("failed to register IB execution command task")?;
352
353 Ok(())
354 }
355
356 fn cached_order_for_modify(&self, client_order_id: &ClientOrderId) -> Option<OrderAny> {
357 self.core.cache().order(client_order_id).map(|o| o.clone())
358 }
359
360 fn reserve_next_local_order_id(next_order_id: &Arc<Mutex<i32>>) -> anyhow::Result<i32> {
361 let mut guard = next_order_id.lock();
362 anyhow::ensure!(
363 *guard > 0,
364 "No valid Interactive Brokers order ID available"
365 );
366 let order_id = *guard;
367 *guard += 1;
368 Ok(order_id)
369 }
370
371 fn apply_client_order_id_floor(next_id: i32, client_id: i32) -> i32 {
372 let client_slot = client_id.unsigned_abs() % 1000;
373 if client_slot == 0 {
374 return next_id;
375 }
376
377 let order_id_floor = (client_slot as i32) * 1_000_000;
378 if next_id > order_id_floor {
379 next_id
380 } else {
381 order_id_floor.saturating_add(next_id.max(1))
382 }
383 }
384
385 async fn get_next_order_id(&self) -> anyhow::Result<i32> {
391 let client = self.ib_client.as_ref().context("IB client not connected")?;
392
393 let timeout_dur = Duration::from_secs(self.config.request_timeout);
394 let order_id = tokio::time::timeout(timeout_dur, client.next_valid_order_id())
395 .await
396 .context("Timeout getting next order ID")??;
397 Ok(order_id)
398 }
399
400 async fn get_highest_open_order_id(&self, client: &Client) -> anyhow::Result<Option<i32>> {
401 let timeout_dur = Duration::from_secs(self.config.request_timeout);
402 let subscription = tokio::time::timeout(timeout_dur, client.all_open_orders())
403 .await
404 .context("Timeout requesting open orders for next order ID initialization")??;
405 let mut subscription = subscription.filter_data();
406 let mut highest_order_id = None;
407
408 while let Some(order_result) = subscription.next().await {
409 match order_result {
410 Ok(Orders::OrderData(data)) => {
411 highest_order_id = Some(
412 highest_order_id
413 .map_or(data.order_id, |current: i32| current.max(data.order_id)),
414 );
415 }
416 Ok(_) => {}
417 Err(e) => {
418 tracing::debug!(
419 "Ignoring open-order event while initializing next order ID: {e}"
420 );
421 }
422 }
423 }
424
425 Ok(highest_order_id)
426 }
427
428 fn begin_task_shutdown(&self) {
429 self.pending_tasks.begin_shutdown();
430 self.session_tasks.begin_shutdown();
431 self.is_connected.store(false, Ordering::Release);
432 self.core.set_disconnected();
433 }
434
435 async fn finish_tasks(&self) -> anyhow::Result<()> {
436 let (session_result, pending_result) = tokio::join!(
437 self.session_tasks
438 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
439 self.pending_tasks
440 .finish_shutdown(Duration::from_secs(2), Duration::from_secs(2)),
441 );
442 session_result.context("failed to finish IB execution session tasks")?;
443 pending_result.context("failed to finish IB execution command tasks")?;
444 Ok(())
445 }
446
447 async fn prepare_task_groups(&mut self) -> anyhow::Result<()> {
448 if !self.session_tasks.is_open() || !self.pending_tasks.is_open() {
449 self.begin_task_shutdown();
450 self.finish_tasks().await?;
451 self.ib_client = None;
452 self.session_tasks
453 .start_generation()
454 .context("failed to start IB execution session task generation")?;
455 self.pending_tasks
456 .start_generation()
457 .context("failed to start IB execution command task generation")?;
458 }
459 Ok(())
460 }
461
462 async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
463 self.begin_task_shutdown();
464 self.ib_client = None;
465 let tasks_result = self.finish_tasks().await;
466 self.is_connected.store(false, Ordering::Release);
467 self.core.set_disconnected();
468 tasks_result
469 }
470}
471
472#[async_trait::async_trait(?Send)]
474impl ExecutionClient for InteractiveBrokersExecutionClient {
475 fn is_connected(&self) -> bool {
476 self.is_connected.load(Ordering::Relaxed)
477 }
478
479 fn client_id(&self) -> ClientId {
480 self.core.client_id
481 }
482
483 fn account_id(&self) -> AccountId {
484 self.core.account_id
485 }
486
487 fn venue(&self) -> Venue {
488 self.core.venue
489 }
490
491 fn handles_order_venue(&self, _venue: Venue) -> bool {
494 true
495 }
496
497 fn oms_type(&self) -> OmsType {
498 self.core.oms_type
499 }
500
501 fn get_account(&self) -> Option<AccountAny> {
502 self.core.cache().account_owned(&self.core.account_id)
503 }
504
505 fn generate_account_state(
506 &self,
507 balances: Vec<AccountBalance>,
508 margins: Vec<MarginBalance>,
509 reported: bool,
510 ts_event: UnixNanos,
511 info: Option<Params>,
512 ) -> anyhow::Result<()> {
513 let factory = OrderEventFactory::new(
514 self.core.trader_id,
515 self.core.account_id,
516 self.core.account_type,
517 self.core.base_currency,
518 );
519 let state = factory.generate_account_state(
520 balances,
521 margins,
522 reported,
523 ts_event,
524 get_atomic_clock_realtime().get_time_ns(),
525 info,
526 );
527 get_exec_event_sender()
528 .send(ExecutionEvent::Account(state))
529 .map_err(|e| anyhow::anyhow!("Failed to send account state: {e}"))
530 }
531
532 fn start(&mut self) -> anyhow::Result<()> {
533 Ok(())
535 }
536
537 fn stop(&mut self) -> anyhow::Result<()> {
538 self.begin_task_shutdown();
539 Ok(())
540 }
541
542 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
543 let order = self.core.get_order(&cmd.client_order_id)?;
544 if let Err(reason) = validate_order(&order) {
545 let reason = reason.to_string();
546 Self::send_order_denied(
547 cmd.order_init.trader_id,
548 cmd.strategy_id,
549 cmd.instrument_id,
550 cmd.order_init.client_order_id,
551 &reason,
552 )?;
553 return Ok(());
554 }
555
556 if let Err(reason) = self.ensure_client_ready_for_order_request("submit order") {
557 self.deny_submit_order_not_ready(&cmd, &reason)?;
558 return Ok(());
559 }
560
561 let client = self.ib_client.as_ref().context("IB client not connected")?;
562
563 let order_id_map = Arc::clone(&self.order_id_map);
564 let venue_order_id_map = Arc::clone(&self.venue_order_id_map);
565 let instrument_id_map = Arc::clone(&self.instrument_id_map);
566 let trader_id_map = Arc::clone(&self.trader_id_map);
567 let strategy_id_map = Arc::clone(&self.strategy_id_map);
568 let active_order_contexts = Arc::clone(&self.active_order_contexts);
569 let terminal_order_contexts = Arc::clone(&self.terminal_order_contexts);
570 let next_order_id = Arc::clone(&self.next_order_id);
571 let instrument_provider = Arc::clone(&self.instrument_provider);
572 let exec_sender = get_exec_event_sender();
573 let clock = get_atomic_clock_realtime();
574 let order_submit_lock = Arc::clone(&self.order_submit_lock);
575
576 let client_clone = client.as_arc().clone();
577
578 let account_id = self.core.account_id;
579
580 let future = async move {
581 if let Err(e) = Self::handle_submit_order_async(
582 &cmd,
583 &client_clone,
584 &order_id_map,
585 &venue_order_id_map,
586 &instrument_id_map,
587 &trader_id_map,
588 &strategy_id_map,
589 &active_order_contexts,
590 &terminal_order_contexts,
591 &next_order_id,
592 &instrument_provider,
593 &exec_sender,
594 clock,
595 account_id,
596 &order_submit_lock,
597 )
598 .await
599 {
600 tracing::error!("Error submitting order: {e}");
601 }
602 };
603
604 self.pending_tasks
605 .spawn(future)
606 .context("failed to register IB execution command task")?;
607
608 Ok(())
609 }
610
611 async fn connect(&mut self) -> anyhow::Result<()> {
612 if self.is_connected.load(Ordering::Relaxed)
613 && self.session_tasks.is_open()
614 && self.pending_tasks.is_open()
615 {
616 log::debug!("Interactive Brokers execution client already connected");
617 return Ok(());
618 }
619
620 self.prepare_task_groups().await?;
621 let setup_guard =
622 TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {});
623
624 tracing::info!("Connecting Interactive Brokers execution client...");
625 log::debug!(
626 "Execution client config host={} port={} client_id={} account_id={:?} request_timeout={} connection_timeout={} fetch_all_open_orders={} track_option_exercise_from_position_update={}",
627 self.config.host,
628 self.config.port,
629 self.config.client_id,
630 self.config.account_id,
631 self.config.request_timeout,
632 self.config.connection_timeout,
633 self.config.fetch_all_open_orders,
634 self.config.track_option_exercise_from_position_update
635 );
636
637 let handle = crate::common::shared_client::get_or_connect(
638 &self.config.host,
639 self.config.port,
640 self.config.client_id,
641 self.config.connection_timeout,
642 )
643 .await
644 .context("Failed to connect to IB Gateway/TWS")?;
645 let client = Arc::clone(handle.as_arc());
646
647 tracing::info!(
648 "Connected to IB Gateway/TWS at {}:{} (client_id: {})",
649 self.config.host,
650 self.config.port,
651 self.config.client_id
652 );
653
654 log::debug!("Initializing IB execution instrument provider");
656
657 if let Err(e) = self
658 .instrument_provider
659 .initialize_with_client(client.as_ref())
660 .await
661 {
662 if !self.config.instrument_provider.load_ids.is_empty()
663 || !self.config.instrument_provider.load_contracts.is_empty()
664 {
665 return Err(e).context("Failed to load configured IB instruments on startup");
666 }
667
668 tracing::warn!("Failed to load instruments on startup: {}", e);
669 }
670
671 self.ib_client = Some(handle);
672
673 let session_result = async {
674
675 log::debug!("Preloading cached spread instruments for execution client");
676 self.preload_cached_spread_instruments(client.as_ref())
677 .await?;
678
679 log::debug!("Requesting next valid IB order ID");
681 let next_id = self.get_next_order_id().await?;
682 log::debug!("Requesting highest open IB order ID");
683 let highest_open_order_id = self.get_highest_open_order_id(client.as_ref()).await?;
684 let client_scoped_next_id =
685 Self::apply_client_order_id_floor(next_id, self.config.client_id);
686 let starting_order_id = highest_open_order_id
687 .map(|order_id| next_id.max(order_id.saturating_add(1)))
688 .unwrap_or(next_id)
689 .max(client_scoped_next_id);
690
691 if starting_order_id != next_id {
692 tracing::debug!(
693 "Adjusted next Interactive Brokers order ID from {} to {} based on client ID/open orders",
694 next_id,
695 starting_order_id
696 );
697 } else {
698 tracing::debug!(
699 "Initialized next Interactive Brokers order ID to {}",
700 starting_order_id
701 );
702 }
703 {
704 let mut id = self
705 .next_order_id
706 .lock();
707 *id = starting_order_id;
708 }
709
710 log::debug!("Starting IB order update stream");
712 self.start_order_updates().await?;
713
714 let client_for_account = Arc::clone(&client);
717 let account_id = self.core.account_id;
718 let _exec_client_core = self.core.clone(); log::debug!("Subscribing to IB account summary for {}", account_id);
720 match crate::execution::account::subscribe_account_summary(&client_for_account, account_id)
721 .await
722 {
723 Ok((balances, margins, info)) => {
724 tracing::debug!(
725 "Received account summary: {} balances, {} margins",
726 balances.len(),
727 margins.len()
728 );
729 let ts_event = get_atomic_clock_realtime().get_time_ns();
731
732 if let Err(e) = ExecutionClient::generate_account_state(
733 self, balances, margins, true, ts_event, info,
735 ) {
736 tracing::warn!("Failed to generate account state: {}", e);
737 }
738 }
739 Err(e) => {
740 tracing::warn!("Failed to subscribe to account summary: {}", e);
741 }
742 }
743
744 let client_for_positions_init = Arc::clone(&client);
747 let position_tracker_init = Arc::clone(&self.position_tracker);
748
749 log::debug!("Initializing IB execution position tracking");
750 if let Err(e) = crate::execution::account::initialize_position_tracking(
751 &client_for_positions_init,
752 self.core.account_id,
753 position_tracker_init,
754 )
755 .await
756 {
757 tracing::warn!("Failed to initialize position tracking: {}", e);
758 }
759
760 let client_for_pnl = Arc::clone(&client); log::debug!("Subscribing to IB PnL updates");
764
765 if let Err(e) = crate::execution::account::subscribe_pnl(
766 &client_for_pnl,
767 self.core.account_id,
768 &self.session_tasks,
769 )
770 .await
771 {
772 tracing::warn!("Failed to subscribe to PnL: {}", e);
773 }
774
775 if self.config.track_option_exercise_from_position_update {
777 let client_for_positions = Arc::clone(&client);
778 let position_tracker_clone = Arc::clone(&self.position_tracker);
779 let instrument_provider_clone = Arc::clone(&self.instrument_provider);
780
781 log::debug!("Subscribing to IB position updates for option exercise tracking");
782
783 if let Err(e) = crate::execution::account::subscribe_positions(
784 &client_for_positions,
785 self.core.account_id,
786 position_tracker_clone,
787 instrument_provider_clone,
788 &self.session_tasks,
789 )
790 .await
791 {
792 tracing::warn!("Failed to subscribe to positions: {}", e);
793 }
794 }
795
796 Ok::<(), anyhow::Error>(())
797 }
798 .await;
799
800 if let Err(e) = session_result {
801 if let Err(teardown_error) = self.teardown_partial_connect().await {
802 return Err(e.context(format!(
803 "IB execution startup teardown failed: {teardown_error}"
804 )));
805 }
806 return Err(e);
807 }
808
809 self.is_connected.store(true, Ordering::Relaxed);
810 self.core.set_connected();
811 setup_guard.disarm();
812
813 tracing::info!("Connected Interactive Brokers execution client");
814 Ok(())
815 }
816
817 async fn disconnect(&mut self) -> anyhow::Result<()> {
818 if !self.is_connected.load(Ordering::Relaxed)
819 && self.ib_client.is_none()
820 && self.session_tasks.is_open()
821 && self.session_tasks.is_empty()
822 && self.pending_tasks.is_open()
823 && self.pending_tasks.is_empty()
824 {
825 log::debug!("Interactive Brokers execution client already disconnected");
826 return Ok(());
827 }
828
829 tracing::info!("Disconnecting Interactive Brokers execution client...");
830
831 self.teardown_partial_connect().await?;
832
833 tracing::info!("Disconnected Interactive Brokers execution client");
834 Ok(())
835 }
836
837 async fn generate_order_status_report(
838 &self,
839 cmd: &GenerateOrderStatusReport,
840 ) -> anyhow::Result<Option<OrderStatusReport>> {
841 let plural_cmd = GenerateOrderStatusReports {
842 command_id: cmd.command_id,
843 ts_init: cmd.ts_init,
844 open_only: false,
845 instrument_id: cmd.instrument_id,
846 start: None,
847 end: None,
848 params: cmd.params.clone(),
849 log_receipt_level: LogLevel::Info,
850 correlation_id: cmd.correlation_id,
851 causation_id: cmd.causation_id,
852 };
853
854 let reports = self.generate_order_status_reports(&plural_cmd).await?;
855
856 let report = reports.into_iter().find(|r| {
858 let matches_client = if let Some(filter_client_id) = cmd.client_order_id {
859 r.client_order_id == Some(filter_client_id)
860 } else {
861 true
862 };
863 let matches_venue = if let Some(filter_venue_id) = cmd.venue_order_id {
864 r.venue_order_id == filter_venue_id
865 } else {
866 true
867 };
868 matches_client && matches_venue
869 });
870
871 Ok(report)
872 }
873
874 async fn generate_order_status_reports(
875 &self,
876 cmd: &GenerateOrderStatusReports,
877 ) -> anyhow::Result<Vec<OrderStatusReport>> {
878 let client = self.ib_client.as_ref().context("IB client not connected")?;
879
880 let timeout_dur = Duration::from_secs(self.config.request_timeout);
881 let subscription = tokio::time::timeout(timeout_dur, client.all_open_orders())
882 .await
883 .context("Timeout requesting open orders")??;
884 let mut subscription = subscription.filter_data();
885 let mut reports = Vec::new();
886 let mut open_order_fills: AHashMap<InstrumentId, Decimal> = AHashMap::new();
887 let ts_init = get_atomic_clock_realtime().get_time_ns();
888 let raw_account_id = raw_ib_account_code(&self.core.account_id);
889
890 while let Some(order_result) = subscription.next().await {
891 match order_result {
892 Ok(Orders::OrderData(data)) => {
893 if !data.order.account.is_empty() && data.order.account != raw_account_id {
894 continue;
895 }
896
897 let instrument_id =
899 match self.resolve_report_contract_instrument_id(&data.contract) {
900 Ok(instrument_id) => instrument_id,
901 Err(e) => {
902 tracing::warn!(
903 order_id = data.order_id,
904 sec_type = ?data.contract.security_type,
905 symbol = data.contract.symbol.as_str(),
906 con_id = data.contract.contract_id,
907 error = %e,
908 "Failed to resolve IBKR order status report instrument ID",
909 );
910 continue;
911 }
912 };
913
914 if let Some(filter_id) = cmd.instrument_id {
916 if instrument_id != filter_id {
917 continue;
918 }
919 }
920
921 match parse_order_status_to_report(
924 &IBOrderStatus {
925 order_id: data.order_id,
926 status: data.order_state.status,
927 filled: data.order.filled_quantity,
928 remaining: (data.order.total_quantity - data.order.filled_quantity)
929 .max(0.0),
930 average_fill_price: None, perm_id: data.order.perm_id,
932 parent_id: 0, last_fill_price: None, client_id: data.order.client_id,
935 why_held: String::new(), market_cap_price: None, },
938 Some(&data.order),
939 instrument_id,
940 self.core.account_id,
941 &self.instrument_provider,
942 ts_init,
943 ) {
944 Ok(report) => {
945 if !cmd.open_only && report.filled_qty.as_decimal() > Decimal::ZERO {
946 let signed_filled = if report.order_side == Some(OrderSide::Buy) {
947 report.filled_qty.as_decimal()
948 } else {
949 -report.filled_qty.as_decimal()
950 };
951 open_order_fills
952 .entry(report.instrument_id)
953 .and_modify(|qty| *qty += signed_filled)
954 .or_insert(signed_filled);
955 }
956 reports.push(report);
957 }
958 Err(e) => {
959 tracing::warn!("Failed to parse order status report: {e}");
960 }
961 }
962 }
963 Ok(_) => {
964 }
966 Err(e) => {
967 tracing::warn!("Error receiving order data: {e}");
968 }
969 }
970 }
971
972 if !cmd.open_only {
973 let positions = tokio::time::timeout(timeout_dur, client.positions())
974 .await
975 .context("Timeout requesting positions for synthetic order reports")??;
976 let mut positions = positions.filter_data();
977
978 while let Some(position_result) = positions.next().await {
979 match position_result {
980 Ok(PositionUpdate::Position(position)) => {
981 if position.account != raw_account_id {
982 continue;
983 }
984
985 let instrument = match self
986 .instrument_provider
987 .get_instrument(client.as_arc().as_ref(), &position.contract)
988 .await
989 {
990 Ok(Some(instrument)) => instrument,
991 Ok(None) => {
992 tracing::warn!(
993 con_id = position.contract.contract_id,
994 sec_type = ?position.contract.security_type,
995 "Cannot generate synthetic order report: instrument not found",
996 );
997 continue;
998 }
999 Err(e) => {
1000 tracing::warn!(
1001 con_id = position.contract.contract_id,
1002 sec_type = ?position.contract.security_type,
1003 error = %e,
1004 "Failed to resolve instrument for synthetic order report",
1005 );
1006 continue;
1007 }
1008 };
1009
1010 let instrument_id = instrument.id();
1011 if let Some(filter_id) = cmd.instrument_id
1012 && instrument_id != filter_id
1013 {
1014 continue;
1015 }
1016
1017 let position_qty =
1018 Decimal::from_f64_retain(position.position).unwrap_or_default();
1019 let open_fills = open_order_fills
1020 .get(&instrument_id)
1021 .copied()
1022 .unwrap_or_default();
1023 let adjusted_qty = position_qty - open_fills;
1024 if adjusted_qty.is_zero() {
1025 continue;
1026 }
1027
1028 let quantity = Quantity::new(
1029 adjusted_qty.abs().to_f64().unwrap_or_default(),
1030 instrument.size_precision(),
1031 );
1032 let order_side = if adjusted_qty > Decimal::ZERO {
1033 OrderSide::Buy
1034 } else {
1035 OrderSide::Sell
1036 };
1037 let id = instrument_id.to_string();
1038 let mut report = OrderStatusReport::new(
1039 self.core.account_id,
1040 instrument_id,
1041 Some(ClientOrderId::new(id.clone())),
1042 VenueOrderId::new(id),
1043 order_side.into(),
1044 OrderType::Market,
1045 TimeInForce::Fok,
1046 OrderStatus::Filled,
1047 quantity,
1048 quantity,
1049 ts_init,
1050 ts_init,
1051 ts_init,
1052 Some(UUID4::new()),
1053 );
1054 report.avg_px = self.position_avg_px_open(
1055 &instrument_id,
1056 &instrument,
1057 position.average_cost,
1058 );
1059 reports.push(report);
1060 }
1061 Ok(PositionUpdate::PositionEnd) => break,
1062 Err(e) => tracing::warn!(
1063 "Error receiving position data for synthetic order report: {e}"
1064 ),
1065 }
1066 }
1067 }
1068
1069 Ok(reports)
1070 }
1071
1072 async fn generate_fill_reports(
1073 &self,
1074 cmd: GenerateFillReports,
1075 ) -> anyhow::Result<Vec<FillReport>> {
1076 let client = self.ib_client.as_ref().context("IB client not connected")?;
1077
1078 let account_code = self.core.account_id.to_string();
1080
1081 let time_filter = if let Some(start) = cmd.start {
1083 let start_dt = start.to_datetime_utc();
1084 start_dt.strftime("%Y%m%d-%H:%M:%S").to_string()
1085 } else {
1086 String::new()
1087 };
1088
1089 let filter = ExecutionFilter {
1090 client_id: None,
1091 account_code,
1092 time: time_filter,
1093 symbol: String::new(),
1094 security_type: String::new(),
1095 exchange: String::new(),
1096 side: None,
1097 last_n_days: 0,
1098 specific_dates: Vec::new(),
1099 };
1100
1101 let timeout_dur = Duration::from_secs(self.config.request_timeout);
1102 let subscription = tokio::time::timeout(timeout_dur, client.executions(filter))
1103 .await
1104 .context("Timeout requesting executions")??;
1105 let mut subscription = subscription.filter_data();
1106 let mut reports = Vec::new();
1107 let ts_init = get_atomic_clock_realtime().get_time_ns();
1108 let mut pending_exec_data: AHashMap<String, ExecutionData> = AHashMap::new();
1109 let mut pending_commissions: AHashMap<String, (f64, String)> = AHashMap::new();
1110
1111 while let Some(exec_result) = subscription.next().await {
1112 match exec_result {
1113 Ok(Executions::ExecutionData(exec_data)) => {
1114 let execution_id = exec_data.execution.execution_id.clone();
1115 if let Some((commission, commission_currency)) =
1116 pending_commissions.remove(&execution_id)
1117 {
1118 if let Some(report) = self.parse_historical_fill_report(
1119 &cmd,
1120 &exec_data,
1121 commission,
1122 &commission_currency,
1123 ts_init,
1124 ) {
1125 reports.push(report);
1126 }
1127 } else {
1128 pending_exec_data.insert(execution_id, exec_data);
1129 }
1130 }
1131 Ok(Executions::CommissionReport(commission)) => {
1132 if let Some(exec_data) = pending_exec_data.remove(&commission.execution_id) {
1133 if let Some(report) = self.parse_historical_fill_report(
1134 &cmd,
1135 &exec_data,
1136 commission.commission,
1137 &commission.currency,
1138 ts_init,
1139 ) {
1140 reports.push(report);
1141 }
1142 } else {
1143 pending_commissions.insert(
1144 commission.execution_id,
1145 (commission.commission, commission.currency),
1146 );
1147 }
1148 }
1149 Err(e) => {
1150 tracing::warn!("Error receiving execution data: {e}");
1151 }
1152 }
1153 }
1154
1155 if !pending_exec_data.is_empty() {
1156 tracing::warn!(
1157 "Skipped {} historical fill reports because IB did not provide matching commission reports",
1158 pending_exec_data.len()
1159 );
1160 }
1161
1162 Ok(reports)
1163 }
1164
1165 async fn generate_position_status_reports(
1166 &self,
1167 cmd: &GeneratePositionStatusReports,
1168 ) -> anyhow::Result<Vec<PositionStatusReport>> {
1169 let client = self.ib_client.as_ref().context("IB client not connected")?;
1170
1171 let timeout_dur = Duration::from_secs(self.config.request_timeout);
1172 let subscription = tokio::time::timeout(timeout_dur, client.positions())
1173 .await
1174 .context("Timeout requesting positions")??;
1175 let mut subscription = subscription.filter_data();
1176 let mut reports = Vec::new();
1177 let ts_init = get_atomic_clock_realtime().get_time_ns();
1178 let raw_account_id = raw_ib_account_code(&self.core.account_id);
1179
1180 while let Some(position_result) = subscription.next().await {
1183 match position_result {
1184 Ok(PositionUpdate::Position(position)) => {
1185 if position.account != raw_account_id {
1187 continue;
1188 }
1189
1190 let instrument = match self
1191 .instrument_provider
1192 .get_instrument(client.as_arc().as_ref(), &position.contract)
1193 .await
1194 {
1195 Ok(Some(instrument)) => instrument,
1196 Ok(None) => {
1197 tracing::warn!(
1198 con_id = position.contract.contract_id,
1199 sec_type = ?position.contract.security_type,
1200 "Cannot generate position status report: instrument not found",
1201 );
1202 continue;
1203 }
1204 Err(e) => {
1205 tracing::warn!(
1206 con_id = position.contract.contract_id,
1207 sec_type = ?position.contract.security_type,
1208 error = %e,
1209 "Failed to resolve position instrument",
1210 );
1211 continue;
1212 }
1213 };
1214 let instrument_id = instrument.id();
1215
1216 if let Some(filter_id) = cmd.instrument_id
1218 && instrument_id != filter_id
1219 {
1220 continue;
1221 }
1222
1223 let position_side = if position.position == 0.0 {
1225 PositionSide::Flat
1226 } else if position.position > 0.0 {
1227 PositionSide::Long
1228 } else {
1229 PositionSide::Short
1230 };
1231
1232 let quantity =
1233 Quantity::new(position.position.abs(), instrument.size_precision());
1234
1235 let avg_px_open = self.position_avg_px_open(
1238 &instrument_id,
1239 &instrument,
1240 position.average_cost,
1241 );
1242
1243 let report = PositionStatusReport::new(
1244 self.core.account_id,
1245 instrument_id,
1246 position_side,
1247 quantity,
1248 ts_init, ts_init, None, None, avg_px_open,
1253 );
1254
1255 reports.push(report);
1256 }
1257 Ok(PositionUpdate::PositionEnd) => {
1258 break;
1260 }
1261 Err(e) => {
1262 tracing::warn!("Error receiving position data: {e}");
1263 }
1264 }
1265 }
1266
1267 if reports.is_empty()
1268 && let Some(instrument_id) = cmd.instrument_id
1269 {
1270 let precision = self
1271 .instrument_provider
1272 .find(&instrument_id)
1273 .map_or(0, |instrument| instrument.size_precision());
1274 reports.push(PositionStatusReport::new(
1275 self.core.account_id,
1276 instrument_id,
1277 PositionSide::Flat,
1278 Quantity::zero(precision),
1279 ts_init,
1280 ts_init,
1281 None,
1282 None,
1283 None,
1284 ));
1285 }
1286
1287 Ok(reports)
1288 }
1289
1290 async fn generate_mass_status(
1291 &self,
1292 lookback_mins: Option<u64>,
1293 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
1294 let ts_now = get_atomic_clock_realtime().get_time_ns();
1295 let start = lookback_mins
1296 .map(DurationNanos::try_from_mins)
1297 .transpose()?
1298 .map(|lookback| ts_now.saturating_sub(lookback));
1299
1300 let order_cmd = GenerateOrderStatusReportsBuilder::default()
1301 .ts_init(ts_now)
1302 .open_only(false)
1303 .start(start)
1304 .build()
1305 .map_err(|e| anyhow::anyhow!("{e}"))?;
1306
1307 let fill_cmd = GenerateFillReportsBuilder::default()
1308 .ts_init(ts_now)
1309 .start(start)
1310 .build()
1311 .map_err(|e| anyhow::anyhow!("{e}"))?;
1312
1313 let position_cmd = GeneratePositionStatusReportsBuilder::default()
1314 .ts_init(ts_now)
1315 .start(start)
1316 .build()
1317 .map_err(|e| anyhow::anyhow!("{e}"))?;
1318
1319 let (order_reports, fill_reports, position_reports) = tokio::try_join!(
1320 self.generate_order_status_reports(&order_cmd),
1321 self.generate_fill_reports(fill_cmd),
1322 self.generate_position_status_reports(&position_cmd),
1323 )?;
1324
1325 tracing::info!(
1326 "generate_mass_status: {} order reports, {} fill reports, {} position reports",
1327 order_reports.len(),
1328 fill_reports.len(),
1329 position_reports.len()
1330 );
1331
1332 let mut mass_status = ExecutionMassStatus::new(
1333 self.core.client_id,
1334 self.core.account_id,
1335 self.core.venue,
1336 ts_now,
1337 Some(UUID4::new()),
1338 );
1339
1340 mass_status.add_order_reports(order_reports);
1341 mass_status.add_fill_reports(fill_reports);
1342 mass_status.add_position_reports(position_reports);
1343
1344 Ok(Some(mass_status))
1345 }
1346
1347 fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
1348 let client = self.ib_client.as_ref().context("IB client not connected")?;
1349
1350 let client_clone = client.as_arc().clone();
1351 let account_id = self.core.account_id;
1352 let account_type = self.core.account_type;
1353 let base_currency = self.core.base_currency;
1354 let clock = get_atomic_clock_realtime();
1355 let request_timeout_secs = self.config.request_timeout;
1356
1357 let future = async move {
1358 let timeout_dur = Duration::from_secs(request_timeout_secs);
1359 let result = tokio::time::timeout(
1360 timeout_dur,
1361 crate::execution::account::subscribe_account_summary(&client_clone, account_id),
1362 )
1363 .await;
1364
1365 match result {
1366 Ok(Ok((balances, margins, info))) => {
1367 let ts_event = clock.get_time_ns();
1368 let ts_now = clock.get_time_ns();
1369
1370 let account_state = AccountState::new(
1371 account_id,
1372 account_type,
1373 balances,
1374 margins,
1375 true,
1376 UUID4::new(),
1377 ts_event,
1378 ts_now,
1379 base_currency,
1380 )
1381 .with_info(info);
1382
1383 let endpoint = MessagingSwitchboard::portfolio_update_account();
1384 send_account_state(endpoint, &account_state);
1385 }
1386 Ok(Err(e)) => {
1387 tracing::error!("Failed to query account state: {e}");
1388 }
1389 Err(_) => {
1390 tracing::error!("Timeout waiting for account summary");
1391 }
1392 }
1393 };
1394
1395 self.pending_tasks
1396 .spawn(future)
1397 .context("failed to register IB execution command task")?;
1398
1399 Ok(())
1400 }
1401
1402 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
1403 let client = self.ib_client.as_ref().context("IB client not connected")?;
1404 let client_order_id = cmd.client_order_id;
1405 let trader_id = cmd.trader_id;
1406 let strategy_id = cmd.strategy_id;
1407 let instrument_id = cmd.instrument_id;
1408
1409 let target_order = if let Some(venue_order_id) = &cmd.venue_order_id {
1410 IbOrderSelector::from_venue_order_id(venue_order_id)?
1411 } else {
1412 let map = self.order_id_map.lock();
1413 IbOrderSelector::OrderId(
1414 *map.get(&cmd.client_order_id)
1415 .context("No venue order id for client_order_id")?,
1416 )
1417 };
1418
1419 let client_clone = client.as_arc().clone();
1420 let instrument_id_map = Arc::clone(&self.instrument_id_map);
1421 let instrument_provider = Arc::clone(&self.instrument_provider);
1422 let account_id = self.core.account_id;
1423 let exec_sender = get_exec_event_sender();
1424 let ts_init = get_atomic_clock_realtime().get_time_ns();
1425 let request_timeout_secs = self.config.request_timeout;
1426 let pending_cancel_orders = Arc::clone(&self.pending_cancel_orders);
1427 let raw_account_id = raw_ib_account_code(&self.core.account_id);
1428
1429 let future = async move {
1430 let timeout_dur = Duration::from_secs(request_timeout_secs);
1431 let subscription =
1432 match tokio::time::timeout(timeout_dur, client_clone.all_open_orders()).await {
1433 Ok(Ok(s)) => s,
1434 Ok(Err(e)) => {
1435 tracing::error!("query_order: failed to request open orders: {e}");
1436 return;
1437 }
1438 Err(_) => {
1439 tracing::error!("query_order: timeout requesting open orders");
1440 return;
1441 }
1442 };
1443 let mut subscription = subscription.filter_data();
1444
1445 while let Some(order_result) = subscription.next().await {
1446 if let Ok(Orders::OrderData(data)) = order_result {
1447 if !data.order.account.is_empty() && data.order.account != raw_account_id {
1448 continue;
1449 }
1450
1451 if !target_order.matches(data.order_id, data.order.perm_id) {
1452 continue;
1453 }
1454
1455 let instrument_id = instrument_id_map.lock().get(&data.order_id).copied();
1456 let instrument_id = match instrument_id {
1457 Some(id) => id,
1458 None => match instrument_provider
1459 .resolve_instrument_id_for_contract(&data.contract)
1460 {
1461 Ok(id) => id,
1462 Err(e) => {
1463 tracing::warn!("query_order: failed to convert contract: {e}");
1464 return;
1465 }
1466 },
1467 };
1468
1469 let report = match parse_order_status_to_report(
1470 &IBOrderStatus {
1471 order_id: data.order_id,
1472 status: data.order_state.status,
1473 filled: data.order.filled_quantity,
1474 remaining: (data.order.total_quantity - data.order.filled_quantity)
1475 .max(0.0),
1476 average_fill_price: None,
1477 perm_id: data.order.perm_id,
1478 parent_id: 0,
1479 last_fill_price: None,
1480 client_id: data.order.client_id,
1481 why_held: String::new(),
1482 market_cap_price: None,
1483 },
1484 Some(&data.order),
1485 instrument_id,
1486 account_id,
1487 &instrument_provider,
1488 ts_init,
1489 ) {
1490 Ok(r) => r,
1491 Err(e) => {
1492 tracing::warn!("query_order: failed to parse order status: {e}");
1493 return;
1494 }
1495 };
1496
1497 if exec_sender
1498 .send(ExecutionEvent::Report(ExecutionReport::Order(Box::new(
1499 report,
1500 ))))
1501 .is_err()
1502 {
1503 tracing::error!("query_order: failed to send order status report");
1504 }
1505 return;
1506 }
1507 }
1508
1509 let was_pending_cancel = pending_cancel_orders.lock().remove(&client_order_id);
1510
1511 if was_pending_cancel {
1512 let event = OrderCanceled::new(
1513 trader_id,
1514 strategy_id,
1515 instrument_id,
1516 client_order_id,
1517 UUID4::new(),
1518 ts_init,
1519 ts_init,
1520 false,
1521 Some(target_order.venue_order_id()),
1522 Some(account_id),
1523 None,
1524 );
1525
1526 if exec_sender
1527 .send(ExecutionEvent::Order(OrderEventAny::Canceled(event)))
1528 .is_err()
1529 {
1530 tracing::error!("query_order: failed to send inferred order canceled event");
1531 } else {
1532 tracing::debug!(
1533 "query_order: inferred cancel for {} from missing open order {}",
1534 client_order_id,
1535 target_order.label()
1536 );
1537 }
1538 return;
1539 }
1540
1541 tracing::debug!(
1542 "query_order: order {} not found in open orders (may be filled or canceled)",
1543 target_order.label()
1544 );
1545 };
1546
1547 self.pending_tasks
1548 .spawn(future)
1549 .context("failed to register IB execution command task")?;
1550
1551 Ok(())
1552 }
1553
1554 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
1555 let orders = self.core.get_orders_for_list(&cmd.order_list)?;
1556 if let Some(reason) = orders.iter().find_map(|order| validate_order(order).err()) {
1557 self.deny_submit_order_list_not_ready(&cmd, &reason.to_string())?;
1558 return Ok(());
1559 }
1560
1561 if let Err(reason) = self.ensure_client_ready_for_order_request("submit order list") {
1562 self.deny_submit_order_list_not_ready(&cmd, &reason)?;
1563 return Ok(());
1564 }
1565
1566 self.submit_order_list_with_orders(cmd, orders)
1567 }
1568
1569 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
1570 if let Err(reason) = self.ensure_client_ready_for_order_request("modify order") {
1571 Self::send_order_modify_rejected(
1572 &cmd,
1573 &reason,
1574 &get_exec_event_sender(),
1575 get_atomic_clock_realtime().get_time_ns(),
1576 self.core.account_id,
1577 )?;
1578 return Ok(());
1579 }
1580
1581 let client = self.ib_client.as_ref().context("IB client not connected")?;
1582
1583 let order_id_map = Arc::clone(&self.order_id_map);
1584 let venue_order_id_map = Arc::clone(&self.venue_order_id_map);
1585 let instrument_id_map = Arc::clone(&self.instrument_id_map);
1586 let instrument_provider = Arc::clone(&self.instrument_provider);
1587 let exec_sender = get_exec_event_sender();
1588 let clock = get_atomic_clock_realtime();
1589 let account_id = self.core.account_id;
1590 let client_clone = client.as_arc().clone();
1591 let request_timeout_secs = self.config.request_timeout;
1592 let original_order = self
1593 .cached_order_for_modify(&cmd.client_order_id)
1594 .map(Arc::new);
1595
1596 if original_order.is_none() {
1597 tracing::debug!(
1598 "Order {} not found in cache for modify; querying IB open orders",
1599 cmd.client_order_id
1600 );
1601 }
1602
1603 let future = async move {
1604 if let Err(e) = Self::handle_modify_order_async(
1605 &cmd,
1606 &client_clone,
1607 &order_id_map,
1608 &venue_order_id_map,
1609 &instrument_id_map,
1610 &instrument_provider,
1611 &exec_sender,
1612 clock,
1613 account_id,
1614 original_order.as_ref(),
1615 request_timeout_secs,
1616 )
1617 .await
1618 {
1619 let reason = format!("Failed to route modify order to IB: {e:#}");
1620
1621 if let Err(send_error) = Self::send_order_modify_rejected(
1622 &cmd,
1623 &reason,
1624 &exec_sender,
1625 clock.get_time_ns(),
1626 account_id,
1627 ) {
1628 tracing::error!("{reason}; failed to emit OrderModifyRejected: {send_error}");
1629 }
1630 }
1631 };
1632
1633 self.pending_tasks
1634 .spawn(future)
1635 .context("failed to register IB execution command task")?;
1636
1637 Ok(())
1638 }
1639
1640 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
1641 let target_order = Arc::new(self.core.get_order(&cmd.client_order_id)?);
1642 let exec_sender = get_exec_event_sender();
1643 let clock = get_atomic_clock_realtime();
1644 let account_id = self.core.account_id;
1645
1646 if let Err(e) = Self::validate_cancel_order_target(&cmd, &target_order) {
1647 let reason = format!("Failed to resolve cancel order target: {e:#}");
1648 Self::send_order_cancel_rejected(
1649 &target_order,
1650 &reason,
1651 &exec_sender,
1652 clock.get_time_ns(),
1653 account_id,
1654 )?;
1655 return Ok(());
1656 }
1657
1658 if let Err(reason) = self.ensure_client_ready_for_order_request("cancel order") {
1659 Self::send_order_cancel_rejected(
1660 &target_order,
1661 &reason,
1662 &exec_sender,
1663 clock.get_time_ns(),
1664 account_id,
1665 )?;
1666 return Ok(());
1667 }
1668
1669 let client = self.ib_client.as_ref().context("IB client not connected")?;
1670
1671 let order_id_map = Arc::clone(&self.order_id_map);
1672 let venue_order_id_map = Arc::clone(&self.venue_order_id_map);
1673 let instrument_id_map = Arc::clone(&self.instrument_id_map);
1674 let trader_id_map = Arc::clone(&self.trader_id_map);
1675 let strategy_id_map = Arc::clone(&self.strategy_id_map);
1676 let pending_cancel_orders = Arc::clone(&self.pending_cancel_orders);
1677 let client_clone = client.as_arc().clone();
1678 let request_timeout_secs = self.config.request_timeout;
1679
1680 let future = async move {
1681 if let Err(e) = Self::handle_cancel_order_async(
1682 &cmd,
1683 &target_order,
1684 &client_clone,
1685 &order_id_map,
1686 &venue_order_id_map,
1687 &instrument_id_map,
1688 &trader_id_map,
1689 &strategy_id_map,
1690 &pending_cancel_orders,
1691 &exec_sender,
1692 clock.get_time_ns(),
1693 account_id,
1694 request_timeout_secs,
1695 )
1696 .await
1697 {
1698 let reason = format!("Failed to route cancel order to IB: {e:#}");
1699
1700 if let Err(send_error) = Self::send_order_cancel_rejected(
1701 &target_order,
1702 &reason,
1703 &exec_sender,
1704 clock.get_time_ns(),
1705 account_id,
1706 ) {
1707 tracing::error!("{reason}; failed to emit OrderCancelRejected: {send_error}");
1708 }
1709 }
1710 };
1711
1712 self.pending_tasks
1713 .spawn(future)
1714 .context("failed to register IB execution command task")?;
1715
1716 Ok(())
1717 }
1718
1719 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
1720 if cmd.order_side.is_some() {
1722 tracing::warn!(
1723 "Interactive Brokers does not support order_side filtering for cancel all orders; \
1724 ignoring order_side={:?} and canceling all orders",
1725 cmd.order_side
1726 );
1727 }
1728
1729 if self
1732 .ensure_client_ready_for_order_request("cancel orders")
1733 .is_err()
1734 {
1735 return Ok(());
1736 }
1737
1738 let client = self.ib_client.as_ref().context("IB client not connected")?;
1739
1740 let orders_to_cancel: Vec<(ClientOrderId, Option<VenueOrderId>)> = {
1743 let cache = self.core.cache();
1744 let mut orders_to_cancel: Vec<(ClientOrderId, Option<VenueOrderId>)> = cache
1745 .orders_open(
1746 None, Some(&cmd.instrument_id), None, None, None, )
1752 .iter()
1753 .map(|order| (order.client_order_id(), order.venue_order_id()))
1754 .collect();
1755
1756 if orders_to_cancel.is_empty() {
1757 let ib_order_ids: Vec<i32> = {
1758 let instrument_id_map = self.instrument_id_map.lock();
1759 instrument_id_map
1760 .iter()
1761 .filter_map(|(order_id, instrument_id)| {
1762 (*instrument_id == cmd.instrument_id).then_some(*order_id)
1763 })
1764 .collect()
1765 };
1766
1767 let venue_map = self.venue_order_id_map.lock();
1768
1769 orders_to_cancel.extend(ib_order_ids.into_iter().filter_map(|ib_order_id| {
1770 venue_map.get(&ib_order_id).copied().map(|client_order_id| {
1771 (
1772 client_order_id,
1773 Some(VenueOrderId::from(ib_order_id.to_string())),
1774 )
1775 })
1776 }));
1777 }
1778
1779 orders_to_cancel.sort_by_key(|(client_order_id, _)| client_order_id.to_string());
1780 orders_to_cancel.dedup_by_key(|(client_order_id, _)| *client_order_id);
1781 orders_to_cancel
1782 };
1783
1784 if orders_to_cancel.is_empty() {
1785 tracing::debug!("No open orders to cancel");
1786 return Ok(());
1787 }
1788
1789 tracing::debug!(
1790 "Canceling {} open order(s) for instrument {}",
1791 orders_to_cancel.len(),
1792 cmd.instrument_id
1793 );
1794
1795 let client_clone = client.as_arc().clone();
1796 let order_id_map = Arc::clone(&self.order_id_map);
1797 let instrument_id_map = Arc::clone(&self.instrument_id_map);
1798 let trader_id_map = Arc::clone(&self.trader_id_map);
1799 let strategy_id_map = Arc::clone(&self.strategy_id_map);
1800 let pending_cancel_orders = Arc::clone(&self.pending_cancel_orders);
1801 let exec_sender = get_exec_event_sender();
1802 let clock = get_atomic_clock_realtime();
1803 let account_id = self.core.account_id;
1804 let request_timeout_secs = self.config.request_timeout;
1805
1806 let future = async move {
1807 if let Err(e) = Self::handle_cancel_all_orders_async(
1808 &client_clone,
1809 &order_id_map,
1810 &instrument_id_map,
1811 &trader_id_map,
1812 &strategy_id_map,
1813 &pending_cancel_orders,
1814 &exec_sender,
1815 clock.get_time_ns(),
1816 account_id,
1817 request_timeout_secs,
1818 orders_to_cancel,
1819 )
1820 .await
1821 {
1822 tracing::error!("Error canceling all orders: {e}");
1823 }
1824 };
1825
1826 self.pending_tasks
1827 .spawn(future)
1828 .context("failed to register IB execution command task")?;
1829
1830 Ok(())
1831 }
1832
1833 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
1834 for cancel_cmd in cmd.cancels {
1836 self.cancel_order(cancel_cmd)?;
1837 }
1838 Ok(())
1839 }
1840}
1841
1842fn validate_order(order: &impl Order) -> Result<(), OrderDeniedReason> {
1843 if order.is_reduce_only() {
1844 return Err(OrderDeniedReason::UnsupportedReduceOnly);
1845 }
1846
1847 Ok(())
1848}
1849
1850impl InteractiveBrokersExecutionClient {
1851 fn is_ready_for_order_request(&self) -> bool {
1852 if !self.is_connected.load(Ordering::Relaxed) {
1853 return false;
1854 }
1855
1856 if !self
1857 .ib_client
1858 .as_ref()
1859 .is_some_and(|client| client.is_connected())
1860 {
1861 return false;
1862 }
1863
1864 *self.next_order_id.lock() > 0
1865 }
1866
1867 fn ensure_client_ready_for_order_request(&self, request: &str) -> Result<(), String> {
1868 if self.is_ready_for_order_request() {
1869 return Ok(());
1870 }
1871
1872 let reason = format!("Interactive Brokers client is not ready; refusing to {request}");
1873 tracing::warn!("{reason}");
1874 Err(reason)
1875 }
1876
1877 fn deny_submit_order_not_ready(&self, cmd: &SubmitOrder, reason: &str) -> anyhow::Result<()> {
1878 Self::send_order_denied(
1879 cmd.order_init.trader_id,
1880 cmd.strategy_id,
1881 cmd.instrument_id,
1882 cmd.order_init.client_order_id,
1883 reason,
1884 )
1885 }
1886
1887 fn deny_submit_order_list_not_ready(
1888 &self,
1889 cmd: &SubmitOrderList,
1890 reason: &str,
1891 ) -> anyhow::Result<()> {
1892 for order_init in &cmd.order_inits {
1893 Self::send_order_denied(
1894 order_init.trader_id,
1895 cmd.strategy_id,
1896 cmd.instrument_id,
1897 order_init.client_order_id,
1898 reason,
1899 )?;
1900 }
1901
1902 Ok(())
1903 }
1904
1905 fn send_order_denied(
1906 trader_id: TraderId,
1907 strategy_id: StrategyId,
1908 instrument_id: InstrumentId,
1909 client_order_id: ClientOrderId,
1910 reason: &str,
1911 ) -> anyhow::Result<()> {
1912 let ts_event = get_atomic_clock_realtime().get_time_ns();
1913 let event = OrderDenied::new(
1914 trader_id,
1915 strategy_id,
1916 instrument_id,
1917 client_order_id,
1918 Ustr::from(reason),
1919 UUID4::new(),
1920 ts_event,
1921 ts_event,
1922 );
1923
1924 get_exec_event_sender()
1925 .send(ExecutionEvent::Order(OrderEventAny::Denied(event)))
1926 .map_err(|e| anyhow::anyhow!("Failed to send order denied event: {e}"))
1927 }
1928
1929 fn send_order_modify_rejected(
1930 cmd: &ModifyOrder,
1931 reason: &str,
1932 exec_sender: &EventSender<ExecutionEvent>,
1933 ts_event: UnixNanos,
1934 account_id: AccountId,
1935 ) -> anyhow::Result<()> {
1936 let event = OrderModifyRejected::new(
1937 cmd.trader_id,
1938 cmd.strategy_id,
1939 cmd.instrument_id,
1940 cmd.client_order_id,
1941 Ustr::from(reason),
1942 UUID4::new(),
1943 ts_event,
1944 ts_event,
1945 false,
1946 cmd.venue_order_id,
1947 Some(account_id),
1948 );
1949 exec_sender
1950 .send(ExecutionEvent::Order(OrderEventAny::ModifyRejected(event)))
1951 .map_err(|e| anyhow::anyhow!("Failed to send order modify rejected event: {e}"))
1952 }
1953
1954 fn send_order_cancel_rejected(
1955 target_order: &OrderAny,
1956 reason: &str,
1957 exec_sender: &EventSender<ExecutionEvent>,
1958 ts_event: UnixNanos,
1959 account_id: AccountId,
1960 ) -> anyhow::Result<()> {
1961 let event = OrderCancelRejected::new(
1962 target_order.trader_id(),
1963 target_order.strategy_id(),
1964 target_order.instrument_id(),
1965 target_order.client_order_id(),
1966 Ustr::from(reason),
1967 UUID4::new(),
1968 ts_event,
1969 ts_event,
1970 false,
1971 target_order.venue_order_id(),
1972 Some(account_id),
1973 );
1974 exec_sender
1975 .send(ExecutionEvent::Order(OrderEventAny::CancelRejected(event)))
1976 .map_err(|e| anyhow::anyhow!("Failed to send order cancel rejected event: {e}"))
1977 }
1978}
1979
1980#[allow(dead_code)]
1981impl InteractiveBrokersExecutionClient {
1982 fn parse_historical_fill_report(
1983 &self,
1984 cmd: &GenerateFillReports,
1985 exec_data: &ExecutionData,
1986 commission: f64,
1987 commission_currency: &str,
1988 ts_init: UnixNanos,
1989 ) -> Option<FillReport> {
1990 let instrument_id = match self.resolve_historical_execution_instrument_id(exec_data) {
1991 Ok(instrument_id) => instrument_id,
1992 Err(e) => {
1993 Self::warn_historical_fill_report_parse_error(exec_data, &e);
1994 return None;
1995 }
1996 };
1997
1998 if let Some(filter_id) = cmd.instrument_id
1999 && instrument_id != filter_id
2000 {
2001 return None;
2002 }
2003
2004 if let Some(filter_venue_order_id) = cmd.venue_order_id
2005 && ib_venue_order_id(exec_data.execution.order_id, exec_data.execution.perm_id)
2006 != filter_venue_order_id
2007 {
2008 return None;
2009 }
2010
2011 if let Some(end) = cmd.end {
2012 match parse_execution_time(&exec_data.execution.time) {
2013 Ok(ts_event) if ts_event > end => return None,
2014 Ok(_) => {}
2015 Err(e) => {
2016 Self::warn_historical_fill_report_parse_error(exec_data, &e);
2017 return None;
2018 }
2019 }
2020 }
2021
2022 match parse_execution_to_fill_report(
2023 &exec_data.execution,
2024 &exec_data.contract,
2025 commission,
2026 commission_currency,
2027 instrument_id,
2028 self.core.account_id,
2029 &self.instrument_provider,
2030 ts_init,
2031 None, ) {
2033 Ok(report) => Some(report),
2034 Err(e) => {
2035 Self::warn_historical_fill_report_parse_error(exec_data, &e);
2036 None
2037 }
2038 }
2039 }
2040
2041 fn resolve_historical_execution_instrument_id(
2042 &self,
2043 exec_data: &ExecutionData,
2044 ) -> anyhow::Result<InstrumentId> {
2045 self.resolve_report_contract_instrument_id(&exec_data.contract)
2046 }
2047
2048 fn resolve_report_contract_instrument_id(
2049 &self,
2050 contract: &Contract,
2051 ) -> anyhow::Result<InstrumentId> {
2052 match self
2053 .instrument_provider
2054 .resolve_instrument_id_for_contract(contract)
2055 {
2056 Ok(instrument_id) => Ok(instrument_id),
2057 Err(provider_error) if contract.security_type != SecurityType::Spread => {
2058 ib_contract_to_instrument_id_simple(contract).with_context(|| {
2059 format!(
2060 "Failed to resolve IBKR contract to instrument ID using provider ({provider_error}) or simple conversion",
2061 )
2062 })
2063 }
2064 Err(provider_error) => Err(provider_error)
2065 .context("Failed to resolve BAG contract to spread instrument ID"),
2066 }
2067 }
2068
2069 fn position_avg_px_open(
2070 &self,
2071 instrument_id: &InstrumentId,
2072 instrument: &InstrumentAny,
2073 average_cost: f64,
2074 ) -> Option<Decimal> {
2075 if average_cost <= 0.0 {
2076 return None;
2077 }
2078
2079 let price_magnifier = self.instrument_provider.get_price_magnifier(instrument_id) as f64;
2080 let multiplier = instrument.multiplier().as_f64();
2081 let converted_avg_cost = average_cost / (multiplier * price_magnifier);
2082 Decimal::from_f64_retain(converted_avg_cost)
2083 .map(|price| price.round_dp(instrument.price_precision() as u32))
2084 }
2085
2086 fn warn_historical_fill_report_parse_error(exec_data: &ExecutionData, error: &anyhow::Error) {
2087 tracing::warn!(
2088 symbol = exec_data.contract.symbol.as_str(),
2089 sec_type = ?exec_data.contract.security_type,
2090 exchange = exec_data.contract.exchange.as_str(),
2091 primary_exchange = exec_data.contract.primary_exchange.as_str(),
2092 local_symbol = exec_data.contract.local_symbol.as_str(),
2093 con_id = exec_data.contract.contract_id,
2094 order_id = exec_data.execution.order_id,
2095 order_ref = exec_data.execution.order_reference.as_str(),
2096 execution_id = exec_data.execution.execution_id.as_str(),
2097 error = %error,
2098 "Failed to parse IBKR historical fill report",
2099 );
2100 }
2101
2102 fn validate_cancel_order_target(
2103 cmd: &CancelOrder,
2104 target_order: &OrderAny,
2105 ) -> anyhow::Result<()> {
2106 anyhow::ensure!(
2107 cmd.client_order_id == target_order.client_order_id(),
2108 "command client order ID {} does not match cached order {}",
2109 cmd.client_order_id,
2110 target_order.client_order_id()
2111 );
2112 anyhow::ensure!(
2113 cmd.instrument_id == target_order.instrument_id(),
2114 "command instrument ID {} does not match cached order {}",
2115 cmd.instrument_id,
2116 target_order.instrument_id()
2117 );
2118
2119 if let (Some(command_venue_order_id), Some(target_venue_order_id)) =
2121 (cmd.venue_order_id.as_ref(), target_order.venue_order_id())
2122 {
2123 anyhow::ensure!(
2124 command_venue_order_id == &target_venue_order_id,
2125 "command venue order ID {command_venue_order_id} does not match cached order {target_venue_order_id}"
2126 );
2127 }
2128
2129 Ok(())
2130 }
2131
2132 async fn handle_cancel_order_async(
2138 cmd: &CancelOrder,
2139 target_order: &OrderAny,
2140 client: &Arc<Client>,
2141 order_id_map: &Arc<Mutex<AHashMap<ClientOrderId, i32>>>,
2142 venue_order_id_map: &Arc<Mutex<AHashMap<i32, ClientOrderId>>>,
2143 instrument_id_map: &Arc<Mutex<AHashMap<i32, InstrumentId>>>,
2144 trader_id_map: &Arc<Mutex<AHashMap<i32, TraderId>>>,
2145 strategy_id_map: &Arc<Mutex<AHashMap<i32, StrategyId>>>,
2146 pending_cancel_orders: &Arc<Mutex<ahash::AHashSet<ClientOrderId>>>,
2147 exec_sender: &EventSender<ExecutionEvent>,
2148 ts_init: UnixNanos,
2149 account_id: AccountId,
2150 request_timeout_secs: u64,
2151 ) -> anyhow::Result<()> {
2152 let order_selector = if let Some(venue_order_id) = &cmd.venue_order_id {
2153 IbOrderSelector::from_venue_order_id(venue_order_id)?
2154 } else {
2155 let map = order_id_map.lock();
2156 IbOrderSelector::OrderId(
2157 *map.get(&cmd.client_order_id)
2158 .context("No IB order ID mapping found for client order ID")?,
2159 )
2160 };
2161 let ib_order_id =
2162 Self::resolve_ib_order_id(client, order_selector, account_id, request_timeout_secs)
2163 .await?;
2164 Self::cache_cancel_order_tracking(
2165 ib_order_id,
2166 cmd,
2167 target_order,
2168 order_id_map,
2169 venue_order_id_map,
2170 instrument_id_map,
2171 trader_id_map,
2172 strategy_id_map,
2173 )?;
2174
2175 if let Err(e) = client.cancel_order(ib_order_id, "").await {
2176 tracing::error!(
2177 "Cancel outcome is unknown after attempting to send order {} to IB: {e}",
2178 cmd.client_order_id
2179 );
2180 return Ok(());
2181 }
2182
2183 let venue_order_id = target_order
2184 .venue_order_id()
2185 .unwrap_or_else(|| VenueOrderId::from(ib_order_id.to_string()));
2186 if let Err(e) = Self::emit_order_pending_cancel(
2187 ib_order_id,
2188 cmd.client_order_id,
2189 venue_order_id,
2190 instrument_id_map,
2191 trader_id_map,
2192 strategy_id_map,
2193 pending_cancel_orders,
2194 exec_sender,
2195 ts_init,
2196 account_id,
2197 ) {
2198 tracing::error!(
2199 "Cancel request for order {} was sent, but OrderPendingCancel emission failed: {e}",
2200 cmd.client_order_id
2201 );
2202 }
2203
2204 Ok(())
2205 }
2206
2207 async fn resolve_ib_order_id(
2208 client: &Arc<Client>,
2209 order_selector: IbOrderSelector,
2210 account_id: AccountId,
2211 request_timeout_secs: u64,
2212 ) -> anyhow::Result<i32> {
2213 let target_perm_id = match order_selector {
2214 IbOrderSelector::OrderId(order_id) => return Ok(order_id),
2215 IbOrderSelector::PermId(perm_id) => perm_id,
2216 };
2217
2218 let timeout_dur = Duration::from_secs(request_timeout_secs);
2219 let raw_account_id = raw_ib_account_code(&account_id);
2220 let subscription = match tokio::time::timeout(timeout_dur, client.all_open_orders()).await {
2221 Ok(Ok(subscription)) => subscription,
2222 Ok(Err(e)) => anyhow::bail!("Failed to request open orders for perm_id lookup: {e}"),
2223 Err(_) => anyhow::bail!("Timed out requesting open orders for perm_id lookup"),
2224 };
2225 let mut subscription = subscription.filter_data();
2226
2227 while let Some(order_result) = subscription.next().await {
2228 let Orders::OrderData(data) = order_result? else {
2229 continue;
2230 };
2231
2232 if !Self::is_active_open_order(&data.order) {
2233 continue;
2234 }
2235
2236 if !data.order.account.is_empty() && data.order.account != raw_account_id {
2237 continue;
2238 }
2239
2240 if data.order.perm_id != target_perm_id {
2241 continue;
2242 }
2243
2244 if data.order_id == 0 {
2245 anyhow::bail!(
2246 "Cannot resolve PERM-{target_perm_id}: matching open order has no IB order_id"
2247 );
2248 }
2249
2250 return Ok(data.order_id);
2251 }
2252
2253 anyhow::bail!("Cannot resolve PERM-{target_perm_id}: no matching open order found")
2254 }
2255
2256 fn is_active_open_order(order: &ibapi::orders::Order) -> bool {
2257 !order.deactivate
2258 }
2259
2260 fn is_definitive_order_submit_error(error: &ibapi::Error) -> bool {
2261 matches!(
2262 error,
2263 ibapi::Error::InvalidArgument(_) | ibapi::Error::ServerVersion(_, _, _)
2264 )
2265 }
2266
2267 fn classify_order_submit_error(error: &ibapi::Error) -> CommandFailure {
2268 let reason = error.to_string();
2269
2270 if Self::is_definitive_order_submit_error(error) {
2271 CommandFailure::not_sent(reason)
2272 } else if matches!(
2273 error,
2274 ibapi::Error::Notice(notice)
2275 if notice.category() == ibapi::NoticeCategory::OrderRejection
2276 ) {
2277 CommandFailure::venue_rejected(reason)
2278 } else {
2279 CommandFailure::ambiguous(reason)
2280 }
2281 }
2282
2283 async fn handle_cancel_all_orders_async(
2284 client: &Arc<Client>,
2285 order_id_map: &Arc<Mutex<AHashMap<ClientOrderId, i32>>>,
2286 instrument_id_map: &Arc<Mutex<AHashMap<i32, InstrumentId>>>,
2287 trader_id_map: &Arc<Mutex<AHashMap<i32, TraderId>>>,
2288 strategy_id_map: &Arc<Mutex<AHashMap<i32, StrategyId>>>,
2289 pending_cancel_orders: &Arc<Mutex<ahash::AHashSet<ClientOrderId>>>,
2290 exec_sender: &EventSender<ExecutionEvent>,
2291 ts_init: UnixNanos,
2292 account_id: AccountId,
2293 request_timeout_secs: u64,
2294 orders_to_cancel: Vec<(ClientOrderId, Option<VenueOrderId>)>,
2295 ) -> anyhow::Result<()> {
2296 let order_selectors: Vec<(ClientOrderId, IbOrderSelector, Option<VenueOrderId>)> = {
2298 let order_id_map_guard = order_id_map.lock();
2299
2300 orders_to_cancel
2301 .into_iter()
2302 .filter_map(|(client_order_id, venue_order_id)| {
2303 if let Some(venue_order_id) = venue_order_id {
2304 match IbOrderSelector::from_venue_order_id(&venue_order_id) {
2305 Ok(order_selector) => {
2306 return Some((client_order_id, order_selector, Some(venue_order_id)));
2307 }
2308 Err(e) => {
2309 tracing::error!(
2310 "Failed resolve cancel-all order {} from venue order ID {}: {e}",
2311 client_order_id,
2312 venue_order_id
2313 );
2314 return None;
2315 }
2316 }
2317 }
2318
2319 order_id_map_guard
2320 .get(&client_order_id)
2321 .copied()
2322 .map(|ib_order_id| {
2323 (
2324 client_order_id,
2325 IbOrderSelector::OrderId(ib_order_id),
2326 None,
2327 )
2328 })
2329 })
2330 .collect()
2331 };
2332
2333 for (client_order_id, order_selector, venue_order_id) in order_selectors {
2335 let ib_order_id = match Self::resolve_ib_order_id(
2336 client,
2337 order_selector,
2338 account_id,
2339 request_timeout_secs,
2340 )
2341 .await
2342 {
2343 Ok(ib_order_id) => ib_order_id,
2344 Err(e) => {
2345 tracing::error!("Failed resolve cancel-all order {client_order_id}: {e}");
2346 continue;
2347 }
2348 };
2349 let venue_order_id =
2350 venue_order_id.unwrap_or_else(|| VenueOrderId::from(ib_order_id.to_string()));
2351
2352 if let Err(e) = client.cancel_order(ib_order_id, "").await {
2353 tracing::error!(
2354 "Failed to cancel order {} (IB order ID: {}): {e}",
2355 client_order_id,
2356 ib_order_id
2357 );
2358 } else {
2359 if let Err(e) = Self::emit_order_pending_cancel(
2360 ib_order_id,
2361 client_order_id,
2362 venue_order_id,
2363 instrument_id_map,
2364 trader_id_map,
2365 strategy_id_map,
2366 pending_cancel_orders,
2367 exec_sender,
2368 ts_init,
2369 account_id,
2370 ) {
2371 tracing::error!(
2372 "Failed to emit pending cancel for order {} (IB order ID: {}): {e}",
2373 client_order_id,
2374 ib_order_id
2375 );
2376 }
2377 tracing::debug!(
2378 "Canceled order {} (IB order ID: {})",
2379 client_order_id,
2380 ib_order_id
2381 );
2382 }
2383 }
2384
2385 tracing::debug!("Finished canceling all orders");
2386
2387 Ok(())
2388 }
2389
2390 #[allow(clippy::too_many_arguments)]
2391 fn emit_order_pending_cancel(
2392 order_id: i32,
2393 client_order_id: ClientOrderId,
2394 venue_order_id: VenueOrderId,
2395 instrument_id_map: &Arc<Mutex<AHashMap<i32, InstrumentId>>>,
2396 trader_id_map: &Arc<Mutex<AHashMap<i32, TraderId>>>,
2397 strategy_id_map: &Arc<Mutex<AHashMap<i32, StrategyId>>>,
2398 pending_cancel_orders: &Arc<Mutex<ahash::AHashSet<ClientOrderId>>>,
2399 exec_sender: &EventSender<ExecutionEvent>,
2400 ts_init: UnixNanos,
2401 account_id: AccountId,
2402 ) -> anyhow::Result<()> {
2403 let mut pending = pending_cancel_orders.lock();
2404 if !pending.insert(client_order_id) {
2405 return Ok(());
2406 }
2407 drop(pending);
2408
2409 let instrument_id = Self::get_mapped_instrument_id(order_id, instrument_id_map)
2410 .context("Instrument ID not found for pending cancel order")?;
2411 let (trader_id, strategy_id) =
2412 Self::get_required_order_actor_ids(order_id, trader_id_map, strategy_id_map)?;
2413
2414 let event = OrderPendingCancel::new(
2415 trader_id,
2416 strategy_id,
2417 instrument_id,
2418 client_order_id,
2419 Some(account_id),
2420 UUID4::new(),
2421 ts_init,
2422 ts_init,
2423 false,
2424 Some(venue_order_id),
2425 );
2426
2427 exec_sender
2428 .send(ExecutionEvent::Order(OrderEventAny::PendingCancel(event)))
2429 .map_err(|e| anyhow::anyhow!("Failed to send order pending cancel event: {e}"))?;
2430
2431 Ok(())
2432 }
2433}