Skip to main content

nautilus_interactive_brokers/execution/
core.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Core execution client implementation for Interactive Brokers.
17
18#[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/// Interactive Brokers execution client.
120///
121/// This client provides order execution functionality using the `rust-ibapi` library.
122/// It manages order submission, modification, cancellation, and execution reporting.
123#[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    /// Creates a new [`InteractiveBrokersExecutionClient`].
242    ///
243    /// # Arguments
244    ///
245    /// * `core` - Core execution client functionality
246    /// * `config` - Configuration for the client
247    /// * `instrument_provider` - Instrument provider
248    ///
249    /// # Errors
250    ///
251    /// Returns an error if client creation fails.
252    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 account_id is provided in config, use it
264        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    /// Gets the next valid order ID from IB.
386    ///
387    /// # Errors
388    ///
389    /// Returns an error if getting the next order ID fails.
390    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// Implementation of ExecutionClient trait
473#[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    // IB uses a broker venue for the client while routing exchange-MIC instruments;
492    // contract transformation remains the authority for actual venue support.
493    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        // Start is handled by connect() for live clients
534        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        // Initialize provider and load instruments from cache/config if configured
655        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        // Get initial next order ID (uses self.ib_client internally)
680        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        // Start order update subscription (uses self.ib_client internally)
711        log::debug!("Starting IB order update stream");
712        self.start_order_updates().await?;
713
714        // Subscribe to account summary and generate initial account state
715        // Wait for initial account summary to load before proceeding
716        let client_for_account = Arc::clone(&client);
717        let account_id = self.core.account_id;
718        let _exec_client_core = self.core.clone(); // Clone core to generate account state
719        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                // Generate account state event like Python version
730                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, // reported
734                    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        // Initialize position tracking with existing positions
745        // This avoids processing duplicates from execDetails
746        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        // Subscribe to PnL updates
761        let client_for_pnl = Arc::clone(&client); // Clone Arc
762
763        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        // Subscribe to position updates for option exercise tracking if enabled
776        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        // Filter by client_order_id and venue_order_id
857        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                    // Convert IB contract to instrument ID
898                    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                    // Filter by instrument_id if specified
915                    if let Some(filter_id) = cmd.instrument_id {
916                        if instrument_id != filter_id {
917                            continue;
918                        }
919                    }
920
921                    // Parse to order status report using minimal OrderStatus
922                    // Note: OrderState doesn't have filled/average_fill_price, so we use defaults
923                    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, // Not available in OrderState
931                            perm_id: data.order.perm_id,
932                            parent_id: 0,          // Not available in OrderState
933                            last_fill_price: None, // Not available in OrderState
934                            client_id: data.order.client_id,
935                            why_held: String::new(), // Not available in OrderState
936                            market_cap_price: None,  // Not available in OrderState
937                        },
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                    // Ignore other order types
965                }
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        // Get account code from account ID
1079        let account_code = self.core.account_id.to_string();
1080
1081        // Build time filter from start if provided.
1082        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        // Process positions until PositionEnd; return empty list when none (reconciliation parity:
1181        // never return None/missing for "no positions").
1182        while let Some(position_result) = subscription.next().await {
1183            match position_result {
1184                Ok(PositionUpdate::Position(position)) => {
1185                    // Filter for the specific account
1186                    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                    // Filter by instrument_id if specified
1217                    if let Some(filter_id) = cmd.instrument_id
1218                        && instrument_id != filter_id
1219                    {
1220                        continue;
1221                    }
1222
1223                    // Determine position side
1224                    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                    // Convert IB avg_cost to Nautilus Price, accounting for price magnifier and multiplier
1236                    // Python: converted_avg_cost = avg_cost / (multiplier * price_magnifier)
1237                    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_last
1249                        ts_init, // ts_init
1250                        None,    // report_id: auto-generated
1251                        None,    // venue_position_id
1252                        avg_px_open,
1253                    );
1254
1255                    reports.push(report);
1256                }
1257                Ok(PositionUpdate::PositionEnd) => {
1258                    // End of position list
1259                    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        // Warn if order_side is specified (IB doesn't support side filtering)
1721        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        // Not-ready warning already logged; a whole-request failure must not
1730        // fan out per-order rejections.
1731        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        // Get open orders from cache before spawning async task (Rc doesn't work across async boundaries)
1741        // Note: In Rust, instrument_id is always required, so we always filter by it
1742        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,                     // venue
1747                    Some(&cmd.instrument_id), // instrument_id (always filter by it in Rust)
1748                    None,                     // strategy_id
1749                    None,                     // account_id
1750                    None,                     // side (IB doesn't support side filtering)
1751                )
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        // Cancel each order in the batch
1835        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, // avg_px (not available in historical fills)
2032        ) {
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        // Command actor IDs identify the requester and are not ownership evidence
2120        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    /// Handles cancel order asynchronously.
2133    ///
2134    /// # Errors
2135    ///
2136    /// Returns an error if broker order resolution or identity caching fails.
2137    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        // Get all IB order selectors first, then drop the guard before awaiting
2297        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        // Now cancel each order (guard is dropped, so we can await)
2334        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}