use super::*;
#[allow(dead_code)]
impl InteractiveBrokersExecutionClient {
#[allow(clippy::too_many_arguments)]
pub(super) async fn handle_submit_order_async(
cmd: &SubmitOrder,
client: &Arc<Client>,
order_id_map: &Arc<Mutex<AHashMap<ClientOrderId, i32>>>,
venue_order_id_map: &Arc<Mutex<AHashMap<i32, ClientOrderId>>>,
instrument_id_map: &Arc<Mutex<AHashMap<i32, InstrumentId>>>,
trader_id_map: &Arc<Mutex<AHashMap<i32, TraderId>>>,
strategy_id_map: &Arc<Mutex<AHashMap<i32, StrategyId>>>,
active_order_contexts: &Arc<Mutex<AHashMap<i32, TrackedOrderContext>>>,
terminal_order_contexts: &Arc<Mutex<FifoCacheMap<i32, TrackedOrderContext, 10_000>>>,
next_order_id: &Arc<Mutex<i32>>,
instrument_provider: &Arc<InteractiveBrokersInstrumentProvider>,
exec_sender: &tokio::sync::mpsc::UnboundedSender<ExecutionEvent>,
clock: &'static AtomicTime,
account_id: AccountId,
order_submit_lock: &Arc<AsyncMutex<()>>,
) -> anyhow::Result<()> {
if cmd.order_init.post_only {
let ts_event = clock.get_time_ns();
let event = OrderDenied::new(
cmd.order_init.trader_id,
cmd.strategy_id,
cmd.instrument_id,
cmd.order_init.client_order_id,
Ustr::from("`post_only` not supported by Interactive Brokers"),
UUID4::new(),
ts_event,
ts_event,
);
exec_sender
.send(ExecutionEvent::Order(OrderEventAny::Denied(event)))
.map_err(|e| anyhow::anyhow!("Failed to send order denied event: {e}"))?;
anyhow::bail!("`post_only` not supported by Interactive Brokers");
}
let is_inverse = instrument_provider
.find(&cmd.instrument_id)
.map(|instrument| instrument.is_inverse())
.unwrap_or(false);
if cmd.order_init.quote_quantity && !is_inverse {
let ts_event = clock.get_time_ns();
let event = OrderDenied::new(
cmd.order_init.trader_id,
cmd.strategy_id,
cmd.instrument_id,
cmd.order_init.client_order_id,
Ustr::from("UNSUPPORTED_QUOTE_QUANTITY"),
UUID4::new(),
ts_event,
ts_event,
);
exec_sender
.send(ExecutionEvent::Order(OrderEventAny::Denied(event)))
.map_err(|e| anyhow::anyhow!("Failed to send order denied event: {e}"))?;
anyhow::bail!("UNSUPPORTED_QUOTE_QUANTITY");
}
if matches!(
cmd.order_init.order_type,
OrderType::TrailingStopMarket | OrderType::TrailingStopLimit
) && let Some(trailing_offset_type) = cmd.order_init.trailing_offset_type
&& trailing_offset_type != TrailingOffsetType::Price
{
let ts_event = clock.get_time_ns();
let reason = format!(
"`TrailingOffsetType` {:?} is not supported (only PRICE is supported)",
trailing_offset_type
);
let event = OrderDenied::new(
cmd.order_init.trader_id,
cmd.strategy_id,
cmd.instrument_id,
cmd.order_init.client_order_id,
Ustr::from(&reason),
UUID4::new(),
ts_event,
ts_event,
);
exec_sender
.send(ExecutionEvent::Order(OrderEventAny::Denied(event)))
.map_err(|e| anyhow::anyhow!("Failed to send order denied event: {e}"))?;
anyhow::bail!("{}", reason);
}
let contract =
Self::resolve_contract_for_instrument(cmd.instrument_id, instrument_provider)?;
let contract = Self::contract_with_order_exchange_param(contract, cmd.params.as_ref())?;
let order_any = OrderAny::try_from(cmd.order_init.clone())
.context("Failed to construct order from `OrderInitialized`")?;
let order_ref = cmd.order_init.client_order_id.to_string();
let _submit_guard = order_submit_lock.lock().await;
let ib_order_id = Self::reserve_next_local_order_id(next_order_id)?;
let mut ib_order = nautilus_order_to_ib_order(
&order_any,
&contract,
instrument_provider,
ib_order_id,
&order_ref,
)
.context("Failed to transform order")?;
let ib_account = account_id
.to_string()
.split_once('-')
.map_or_else(|| account_id.to_string(), |(_, value)| value.to_string());
ib_order.account = ib_account.clone();
ib_order.clearing_account = ib_account;
Self::cache_order_tracking(
ib_order_id,
cmd.order_init.client_order_id,
cmd.instrument_id,
cmd.order_init.trader_id,
cmd.strategy_id,
cmd.order_init.order_side,
cmd.order_init.order_type,
order_id_map,
venue_order_id_map,
instrument_id_map,
trader_id_map,
strategy_id_map,
active_order_contexts,
terminal_order_contexts,
)?;
let ts_event = clock.get_time_ns();
let event = OrderSubmitted::new(
cmd.order_init.trader_id,
cmd.strategy_id,
cmd.instrument_id,
cmd.order_init.client_order_id,
account_id,
UUID4::new(),
ts_event,
ts_event,
);
exec_sender
.send(ExecutionEvent::Order(OrderEventAny::Submitted(event)))
.map_err(|e| anyhow::anyhow!("Failed to send order submitted event: {e}"))?;
if let Err(e) = client.submit_order(ib_order_id, &contract, &ib_order).await {
return Self::handle_order_submit_failure(
&e,
"Failed to submit order",
ib_order_id,
account_id,
ts_event,
order_id_map,
venue_order_id_map,
instrument_id_map,
trader_id_map,
strategy_id_map,
active_order_contexts,
terminal_order_contexts,
exec_sender,
clock,
);
}
Self::emit_order_accepted_if_needed(
ib_order_id,
VenueOrderId::from(ib_order_id.to_string()),
account_id,
ts_event,
active_order_contexts,
exec_sender,
)?;
tracing::debug!(
"Submitted order {} as IB order ID {}",
cmd.order_init.client_order_id,
ib_order_id
);
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn handle_modify_order_async(
cmd: &ModifyOrder,
client: &Arc<Client>,
order_id_map: &Arc<Mutex<AHashMap<ClientOrderId, i32>>>,
venue_order_id_map: &Arc<Mutex<AHashMap<i32, ClientOrderId>>>,
instrument_id_map: &Arc<Mutex<AHashMap<i32, InstrumentId>>>,
instrument_provider: &Arc<InteractiveBrokersInstrumentProvider>,
_exec_sender: &tokio::sync::mpsc::UnboundedSender<ExecutionEvent>,
_clock: &'static AtomicTime,
account_id: AccountId,
original_order: Option<&Arc<OrderAny>>,
request_timeout_secs: u64,
) -> anyhow::Result<()> {
let target_ib_order_id = Self::target_ib_order_id_for_modify(
cmd,
client,
order_id_map,
account_id,
request_timeout_secs,
)
.await?;
if let Some(original_order) = original_order {
let ib_order_id = target_ib_order_id.context("Order ID not found in mapping")?;
let contract =
Self::resolve_contract_for_instrument(cmd.instrument_id, instrument_provider)?;
let contract = Self::contract_with_order_exchange_param(contract, cmd.params.as_ref())?;
let order_ref = original_order.client_order_id().to_string();
let mut ib_order = nautilus_order_to_ib_order(
original_order,
&contract,
instrument_provider,
ib_order_id,
&order_ref,
)
.context("Failed to transform order to IB order")?;
Self::apply_modify_fields_to_ib_order(cmd, &mut ib_order, instrument_provider);
if let Err(e) = client.submit_order(ib_order_id, &contract, &ib_order).await {
if Self::is_definitive_order_submit_error(&e) {
return Err(e).context("IB rejected the modified order before sending it");
}
tracing::error!(
"Modify outcome is unknown after attempting to send order {} to IB: {e}",
cmd.client_order_id
);
return Ok(());
}
tracing::debug!(
"Modified order {} (IB order ID: {})",
cmd.client_order_id,
ib_order_id
);
return Ok(());
}
Self::handle_modify_open_order_async(
cmd,
client,
target_ib_order_id,
order_id_map,
venue_order_id_map,
instrument_id_map,
instrument_provider,
request_timeout_secs,
)
.await
}
async fn target_ib_order_id_for_modify(
cmd: &ModifyOrder,
client: &Arc<Client>,
order_id_map: &Arc<Mutex<AHashMap<ClientOrderId, i32>>>,
account_id: AccountId,
request_timeout_secs: u64,
) -> anyhow::Result<Option<i32>> {
if let Some(venue_order_id) = &cmd.venue_order_id {
let order_selector = IbOrderSelector::from_venue_order_id(venue_order_id)?;
let order_id =
Self::resolve_ib_order_id(client, order_selector, account_id, request_timeout_secs)
.await?;
return Ok(Some(order_id));
}
let map = order_id_map
.lock()
.map_err(|_| anyhow::anyhow!("Failed to lock order ID map"))?;
Ok(map.get(&cmd.client_order_id).copied())
}
fn apply_modify_fields_to_ib_order(
cmd: &ModifyOrder,
ib_order: &mut ibapi::orders::Order,
instrument_provider: &Arc<InteractiveBrokersInstrumentProvider>,
) {
if let Some(quantity) = cmd.quantity {
ib_order.total_quantity = quantity.as_f64();
}
let price_magnifier = instrument_provider.get_price_magnifier(&cmd.instrument_id) as f64;
if let Some(price) = cmd.price {
ib_order.limit_price = Some(price.as_f64() / price_magnifier);
}
if let Some(trigger_price) = cmd.trigger_price {
let converted_trigger_price = trigger_price.as_f64() / price_magnifier;
if matches!(ib_order.order_type.as_str(), "TRAIL" | "TRAIL LIMIT") {
ib_order.trail_stop_price = Some(converted_trigger_price);
} else {
ib_order.aux_price = Some(converted_trigger_price);
}
}
}
#[allow(clippy::too_many_arguments)]
async fn handle_modify_open_order_async(
cmd: &ModifyOrder,
client: &Arc<Client>,
target_ib_order_id: Option<i32>,
order_id_map: &Arc<Mutex<AHashMap<ClientOrderId, i32>>>,
venue_order_id_map: &Arc<Mutex<AHashMap<i32, ClientOrderId>>>,
instrument_id_map: &Arc<Mutex<AHashMap<i32, InstrumentId>>>,
instrument_provider: &Arc<InteractiveBrokersInstrumentProvider>,
request_timeout_secs: u64,
) -> anyhow::Result<()> {
let timeout_dur = Duration::from_secs(request_timeout_secs);
let subscription = tokio::time::timeout(timeout_dur, client.all_open_orders())
.await
.context("Timeout requesting open orders for modify")??;
let mut subscription = subscription.filter_data();
let client_order_id = cmd.client_order_id.to_string();
while let Some(order_result) = subscription.next().await {
match order_result {
Ok(Orders::OrderData(data)) => {
if !Self::is_active_open_order(&data.order) {
continue;
}
let matches_order_id =
target_ib_order_id.is_some_and(|order_id| data.order_id == order_id);
let matches_order_ref = data.order.order_ref == client_order_id;
if !matches_order_id && !matches_order_ref {
continue;
}
let ib_order_id = data.order_id;
let contract = data.contract;
let contract =
Self::contract_with_order_exchange_param(contract, cmd.params.as_ref())?;
let mut ib_order = data.order;
Self::apply_modify_fields_to_ib_order(cmd, &mut ib_order, instrument_provider);
{
let mut map = order_id_map
.lock()
.map_err(|_| anyhow::anyhow!("Failed to lock order ID map"))?;
map.insert(cmd.client_order_id, ib_order_id);
}
{
let mut map = venue_order_id_map
.lock()
.map_err(|_| anyhow::anyhow!("Failed to lock venue order ID map"))?;
map.insert(ib_order_id, cmd.client_order_id);
}
{
let mut map = instrument_id_map
.lock()
.map_err(|_| anyhow::anyhow!("Failed to lock instrument ID map"))?;
map.insert(ib_order_id, cmd.instrument_id);
}
if let Err(e) = client.submit_order(ib_order_id, &contract, &ib_order).await {
if Self::is_definitive_order_submit_error(&e) {
return Err(e)
.context("IB rejected the modified open order before sending it");
}
tracing::error!(
"Modify outcome is unknown after attempting to send open order {} to IB: {e}",
cmd.client_order_id
);
return Ok(());
}
tracing::debug!(
"Modified open order {} (IB order ID: {}) after cache miss",
cmd.client_order_id,
ib_order_id
);
return Ok(());
}
Ok(_) => {}
Err(e) => {
tracing::warn!("Error receiving open order data for modify: {e}");
}
}
}
anyhow::bail!(
"Order not found for modify in IB open orders: client_order_id={}, venue_order_id={:?}",
cmd.client_order_id,
cmd.venue_order_id,
)
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn handle_submit_order_list_async(
cmd: &SubmitOrderList,
orders: &[OrderAny],
client: &Arc<Client>,
order_id_map: &Arc<Mutex<AHashMap<ClientOrderId, i32>>>,
venue_order_id_map: &Arc<Mutex<AHashMap<i32, ClientOrderId>>>,
instrument_id_map: &Arc<Mutex<AHashMap<i32, InstrumentId>>>,
trader_id_map: &Arc<Mutex<AHashMap<i32, TraderId>>>,
strategy_id_map: &Arc<Mutex<AHashMap<i32, StrategyId>>>,
active_order_contexts: &Arc<Mutex<AHashMap<i32, TrackedOrderContext>>>,
terminal_order_contexts: &Arc<Mutex<FifoCacheMap<i32, TrackedOrderContext, 10_000>>>,
next_order_id: &Arc<Mutex<i32>>,
instrument_provider: &Arc<InteractiveBrokersInstrumentProvider>,
exec_sender: &tokio::sync::mpsc::UnboundedSender<ExecutionEvent>,
clock: &'static AtomicTime,
account_id: AccountId,
strategy_id: StrategyId,
order_submit_lock: &Arc<AsyncMutex<()>>,
) -> anyhow::Result<()> {
let num_orders = orders.len();
anyhow::ensure!(!orders.is_empty(), "Cannot submit an empty order list");
let _submit_guard = order_submit_lock.lock().await;
let ib_account = account_id
.to_string()
.split_once('-')
.map_or_else(|| account_id.to_string(), |(_, value)| value.to_string());
let mut ib_order_ids = AHashMap::with_capacity(num_orders);
for order in orders {
let ib_order_id = Self::reserve_next_local_order_id(next_order_id)?;
ib_order_ids.insert(order.client_order_id(), ib_order_id);
}
for order in orders {
if let Some(parent_order_id) = order.parent_order_id()
&& !ib_order_ids.contains_key(&parent_order_id)
{
let map = order_id_map
.lock()
.map_err(|_| anyhow::anyhow!("Failed to lock order ID map"))?;
anyhow::ensure!(
map.contains_key(&parent_order_id),
"Parent order ID {parent_order_id} not found for order {}",
order.client_order_id(),
);
}
}
for (index, order) in orders.iter().enumerate() {
let is_last = index == num_orders - 1;
let ib_order_id = ib_order_ids[&order.client_order_id()];
let order_contract =
Self::resolve_contract_for_instrument(order.instrument_id(), instrument_provider)?;
let order_contract =
Self::contract_with_order_exchange_param(order_contract, cmd.params.as_ref())?;
let order_ref = order.client_order_id().to_string();
let mut ib_order = nautilus_order_to_ib_order(
order,
&order_contract,
instrument_provider,
ib_order_id,
&order_ref,
)
.context("Failed to transform order")?;
ib_order.account = ib_account.clone();
ib_order.clearing_account = ib_account.clone();
ib_order.transmit = is_last;
if let Some(parent_order_id) = order.parent_order_id() {
let parent_ib_order_id =
ib_order_ids.get(&parent_order_id).copied().or_else(|| {
order_id_map
.lock()
.ok()
.and_then(|map| map.get(&parent_order_id).copied())
});
if let Some(parent_ib_order_id) = parent_ib_order_id {
ib_order.parent_id = parent_ib_order_id;
}
}
Self::cache_order_tracking(
ib_order_id,
order.client_order_id(),
order.instrument_id(),
order.trader_id(),
strategy_id,
order.order_side(),
order.order_type(),
order_id_map,
venue_order_id_map,
instrument_id_map,
trader_id_map,
strategy_id_map,
active_order_contexts,
terminal_order_contexts,
)?;
let ts_event = clock.get_time_ns();
let event = OrderSubmitted::new(
order.trader_id(),
strategy_id,
order.instrument_id(),
order.client_order_id(),
account_id,
UUID4::new(),
ts_event,
ts_event,
);
exec_sender
.send(ExecutionEvent::Order(OrderEventAny::Submitted(event)))
.map_err(|e| anyhow::anyhow!("Failed to send order submitted event: {e}"))?;
if let Err(e) = client
.submit_order(ib_order_id, &order_contract, &ib_order)
.await
{
return Self::handle_order_submit_failure(
&e,
"Failed to submit order from list",
ib_order_id,
account_id,
ts_event,
order_id_map,
venue_order_id_map,
instrument_id_map,
trader_id_map,
strategy_id_map,
active_order_contexts,
terminal_order_contexts,
exec_sender,
clock,
);
}
Self::emit_order_accepted_if_needed(
ib_order_id,
VenueOrderId::from(ib_order_id.to_string()),
account_id,
ts_event,
active_order_contexts,
exec_sender,
)?;
tracing::debug!(
"Submitted order {} from list as IB order ID {}",
order.client_order_id(),
ib_order_id,
);
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub(super) fn handle_order_submit_failure(
error: &ibapi::Error,
failure_prefix: &str,
ib_order_id: i32,
account_id: AccountId,
ts_event: UnixNanos,
order_id_map: &Arc<Mutex<AHashMap<ClientOrderId, i32>>>,
venue_order_id_map: &Arc<Mutex<AHashMap<i32, ClientOrderId>>>,
instrument_id_map: &Arc<Mutex<AHashMap<i32, InstrumentId>>>,
trader_id_map: &Arc<Mutex<AHashMap<i32, TraderId>>>,
strategy_id_map: &Arc<Mutex<AHashMap<i32, StrategyId>>>,
active_order_contexts: &Arc<Mutex<AHashMap<i32, TrackedOrderContext>>>,
terminal_order_contexts: &Arc<Mutex<FifoCacheMap<i32, TrackedOrderContext, 10_000>>>,
exec_sender: &tokio::sync::mpsc::UnboundedSender<ExecutionEvent>,
clock: &'static AtomicTime,
) -> anyhow::Result<()> {
match Self::classify_order_submit_error(error) {
CommandFailure::Ambiguous(reason) => {
anyhow::bail!(
"{failure_prefix}; outcome is unknown after possible transmission: {reason}"
);
}
CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
let context = Self::get_tracked_order_context(
ib_order_id,
active_order_contexts,
terminal_order_contexts,
)?
.with_context(|| format!("Tracked order context not found for {ib_order_id}"))?;
Self::remove_order_tracking(
ib_order_id,
context.client_order_id,
order_id_map,
venue_order_id_map,
instrument_id_map,
trader_id_map,
strategy_id_map,
active_order_contexts,
terminal_order_contexts,
)?;
let reason = format!("{failure_prefix}: {reason}");
let event = OrderRejected::new(
context.trader_id,
context.strategy_id,
context.instrument_id,
context.client_order_id,
account_id,
Ustr::from(&reason),
UUID4::new(),
ts_event,
clock.get_time_ns(),
false,
false,
);
exec_sender
.send(ExecutionEvent::Order(OrderEventAny::Rejected(event)))
.map_err(|e| anyhow::anyhow!("Failed to send order rejected event: {e}"))?;
anyhow::bail!(reason);
}
}
}
}
#[cfg(test)]
mod tests {
use nautilus_model::identifiers::Symbol;
use super::*;
fn modify_trigger_cmd() -> ModifyOrder {
ModifyOrder::new(
TraderId::from("TRADER-001"),
Some(ClientId::from("CLIENT-001")),
StrategyId::from("S-001"),
InstrumentId::new(Symbol::from("AAPL"), Venue::from("NASDAQ")),
ClientOrderId::from("O-001"),
Some(VenueOrderId::from("1")),
None,
None,
Some(Price::from("149.50")),
UUID4::new(),
UnixNanos::default(),
None,
None,
)
}
fn instrument_provider() -> Arc<InteractiveBrokersInstrumentProvider> {
Arc::new(InteractiveBrokersInstrumentProvider::new(
crate::config::InteractiveBrokersInstrumentProviderConfig::default(),
))
}
#[rstest::rstest]
fn modify_trailing_stop_routes_trigger_to_trail_stop_price() {
let mut ib_order = ibapi::orders::Order {
order_type: "TRAIL".to_string(),
aux_price: Some(0.5),
trailing_percent: Some(0.25),
..Default::default()
};
InteractiveBrokersExecutionClient::apply_modify_fields_to_ib_order(
&modify_trigger_cmd(),
&mut ib_order,
&instrument_provider(),
);
assert_eq!(ib_order.aux_price, Some(0.5));
assert_eq!(ib_order.trailing_percent, Some(0.25));
assert_eq!(ib_order.trail_stop_price, Some(149.5));
}
#[rstest::rstest]
fn modify_stop_order_routes_trigger_to_aux_price() {
let mut ib_order = ibapi::orders::Order {
order_type: "STP".to_string(),
..Default::default()
};
InteractiveBrokersExecutionClient::apply_modify_fields_to_ib_order(
&modify_trigger_cmd(),
&mut ib_order,
&instrument_provider(),
);
assert_eq!(ib_order.aux_price, Some(149.5));
assert_eq!(ib_order.trail_stop_price, None);
}
}