use std::{cell::RefCell, fmt::Debug, rc::Rc, str::FromStr, sync::LazyLock, time::Duration};
use indexmap::{IndexMap, IndexSet};
use nautilus_common::{
cache::Cache,
clients::{DEFAULT_POSITION_RECONCILIATION_TOLERANCE, ExecutionClient},
clock::Clock,
enums::{LogColor, LogLevel},
live::dst,
log_info,
messages::{
ExecutionReport,
execution::{
QueryOrder, TradingCommand,
report::{
GenerateOrderStatusReport, GenerateOrderStatusReports,
GeneratePositionStatusReports,
},
},
},
msgbus::{self, MessagingSwitchboard, switchboard},
};
use nautilus_core::{
UUID4, UnixNanos,
datetime::{mins_to_nanos, mins_to_secs},
};
use nautilus_execution::{
engine::ExecutionEngine,
reconciliation::{
calculate_reconciliation_price, create_inferred_fill_for_qty,
create_position_reconciliation_venue_order_id, create_reconciliation_rejected,
create_reconciliation_triggered, generate_external_order_status_events,
generate_reconciliation_order_pre_fill_events,
generate_reconciliation_order_snapshot_events, process_mass_status_for_reconciliation,
reconcile_order_report, should_reconciliation_update,
},
};
use nautilus_model::{
enums::{OmsType, OrderSide, OrderStatus, OrderType, TimeInForce},
events::{OrderCanceled, OrderEventAny, OrderFilled, OrderInitialized},
identifiers::{
AccountId, ClientId, ClientOrderId, InstrumentId, PositionId, StrategyId, TradeId,
TraderId, VenueOrderId,
},
instruments::{Instrument, InstrumentAny},
orders::{Order, OrderAny, TRIGGERABLE_ORDER_TYPES},
position::Position,
reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
types::{Price, Quantity},
};
use rust_decimal::Decimal;
use ustr::Ustr;
use super::recency::RecencyMap;
static TAG_VENUE: LazyLock<Ustr> = LazyLock::new(|| Ustr::from("VENUE"));
static TAG_RECONCILIATION: LazyLock<Ustr> = LazyLock::new(|| Ustr::from("RECONCILIATION"));
pub type InstrumentAccountKey = (InstrumentId, AccountId);
type FillKey = (AccountId, InstrumentId, TradeId);
#[expect(clippy::too_many_arguments)]
fn build_cross_zero_leg_report(
instrument: &InstrumentAny,
account_id: AccountId,
instrument_id: InstrumentId,
order_side: OrderSide,
quantity: Decimal,
avg_px: Decimal,
tag: &str,
ts_now: UnixNanos,
venue_ts_last: UnixNanos,
) -> Option<OrderStatusReport> {
let order_qty = Quantity::from_decimal_dp(quantity, instrument.size_precision()).ok()?;
let fill_price = Price::from_decimal_dp(avg_px, instrument.price_precision()).ok();
let venue_order_id = create_position_reconciliation_venue_order_id(
account_id,
instrument_id,
order_side,
OrderType::Market,
order_qty,
fill_price,
None,
Some(tag),
venue_ts_last,
);
let report = OrderStatusReport::new(
account_id,
instrument_id,
None,
venue_order_id,
order_side,
OrderType::Market,
TimeInForce::Gtc,
OrderStatus::Filled,
order_qty,
order_qty,
ts_now,
ts_now,
ts_now,
None,
)
.with_avg_px(avg_px);
Some(report)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum ReportClientCoverage {
Resolved(IndexSet<ClientId>),
Unresolved,
}
#[derive(Debug, Clone)]
pub struct ExternalOrderMetadata {
pub client_order_id: ClientOrderId,
pub venue_order_id: VenueOrderId,
pub instrument_id: InstrumentId,
pub strategy_id: StrategyId,
pub ts_init: UnixNanos,
}
#[derive(Debug, Default)]
pub struct ReconciliationResult {
pub events: Vec<OrderEventAny>,
pub external_orders: Vec<ExternalOrderMetadata>,
}
#[derive(Debug, Default)]
pub struct InflightCheckResult {
pub events: Vec<OrderEventAny>,
pub queries: Vec<TradingCommand>,
}
#[derive(Debug, Default)]
pub(crate) struct OpenOrderReconciliationResult {
pub events: Vec<OrderEventAny>,
pub targeted_queries: Vec<TargetedOrderQuery>,
}
#[derive(Debug, Clone)]
pub(crate) struct TargetedOrderQuery {
client_order_id: ClientOrderId,
responsible_clients: IndexSet<ClientId>,
command: GenerateOrderStatusReport,
}
impl TargetedOrderQuery {
#[cfg(feature = "node")]
pub(crate) const fn client_order_id(&self) -> ClientOrderId {
self.client_order_id
}
}
#[derive(Debug)]
pub(crate) struct TargetedOrderReportResult {
client_order_id: ClientOrderId,
report: Option<OrderStatusReport>,
coverage_complete: bool,
}
#[derive(Debug, Clone)]
pub(crate) struct OpenOrderReportCheck {
pub command: GenerateOrderStatusReports,
pub filtered_orders: Vec<OrderAny>,
pub client_coverage: IndexMap<ClientOrderId, ReportClientCoverage>,
pub start: Option<UnixNanos>,
}
#[derive(Debug, Clone)]
pub(crate) struct PositionReportCheck {
pub command: GeneratePositionStatusReports,
pub client_coverage: IndexMap<InstrumentAccountKey, ReportClientCoverage>,
pub activity_revisions: IndexMap<InstrumentAccountKey, u64>,
}
struct RetainedFillState {
fill_keys: IndexSet<(AccountId, InstrumentId, TradeId)>,
missing_order_ids: IndexSet<(AccountId, InstrumentId, ClientOrderId)>,
missing_venue_order_ids: IndexSet<(AccountId, InstrumentId, VenueOrderId)>,
netting_lifecycle_starts: IndexMap<(AccountId, InstrumentId, StrategyId), UnixNanos>,
}
#[derive(Default)]
struct ReconciliationFillQueue {
pending_fill_keys: IndexSet<FillKey>,
event_fill_keys: IndexMap<UUID4, FillKey>,
}
impl ReconciliationFillQueue {
fn push(&mut self, events: &mut Vec<OrderEventAny>, event: OrderEventAny, fill_key: FillKey) {
let OrderEventAny::Filled(fill) = &event else {
unreachable!("reported fills always create filled events");
};
self.pending_fill_keys.insert(fill_key);
self.event_fill_keys.insert(fill.event_id, fill_key);
events.push(event);
}
}
#[expect(
clippy::struct_excessive_bools,
reason = "config flags mirror the live execution engine configuration surface"
)]
#[derive(Debug, Clone)]
pub struct ExecutionManagerConfig {
pub trader_id: TraderId,
pub reconciliation: bool,
pub lookback_mins: Option<u64>,
pub reconciliation_instrument_ids: IndexSet<InstrumentId>,
pub filter_unclaimed_external: bool,
pub filter_position_reports: bool,
pub filtered_client_order_ids: IndexSet<ClientOrderId>,
pub generate_missing_orders: bool,
pub inflight_check_interval_ms: u32,
pub inflight_threshold_ms: u64,
pub inflight_max_retries: u32,
pub open_check_interval_secs: Option<f64>,
pub open_check_lookback_mins: Option<u64>,
pub open_check_threshold_ns: u64,
pub open_check_missing_retries: u32,
pub open_check_open_only: bool,
pub max_single_order_queries_per_cycle: u32,
pub single_order_query_delay_ms: u32,
pub position_check_interval_secs: Option<f64>,
pub position_check_lookback_mins: u64,
pub position_check_threshold_ns: u64,
pub position_check_retries: u32,
pub purge_closed_orders_buffer_mins: Option<u32>,
pub purge_closed_positions_buffer_mins: Option<u32>,
pub purge_account_events_lookback_mins: Option<u32>,
pub purge_from_database: bool,
}
impl Default for ExecutionManagerConfig {
fn default() -> Self {
Self {
trader_id: TraderId::default(),
reconciliation: true,
lookback_mins: Some(60),
reconciliation_instrument_ids: IndexSet::new(),
filter_unclaimed_external: false,
filter_position_reports: false,
filtered_client_order_ids: IndexSet::new(),
generate_missing_orders: true,
inflight_check_interval_ms: 2_000,
inflight_threshold_ms: 5_000,
inflight_max_retries: 5,
open_check_interval_secs: None,
open_check_lookback_mins: Some(60),
open_check_threshold_ns: 5_000_000_000,
open_check_missing_retries: 5,
open_check_open_only: true,
max_single_order_queries_per_cycle: 5,
single_order_query_delay_ms: 100,
position_check_interval_secs: None,
position_check_lookback_mins: 60,
position_check_threshold_ns: 60_000_000_000,
position_check_retries: 3,
purge_closed_orders_buffer_mins: None,
purge_closed_positions_buffer_mins: None,
purge_account_events_lookback_mins: None,
purge_from_database: false,
}
}
}
impl ExecutionManagerConfig {
#[must_use]
pub fn with_trader_id(mut self, trader_id: TraderId) -> Self {
self.trader_id = trader_id;
self
}
}
#[derive(Debug, Clone)]
struct InflightCheck {
#[allow(dead_code)]
pub client_order_id: ClientOrderId,
pub submitted_at: dst::time::Instant,
pub retry_count: u32,
pub last_query_at: Option<dst::time::Instant>,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum PositionReportShape {
Unambiguous,
MultiLeg,
}
#[derive(Clone, Copy)]
struct PositionReconciliationState {
report_shape: PositionReportShape,
retries: u32,
}
#[derive(Clone)]
pub struct ExecutionManager {
clock: Rc<RefCell<dyn Clock>>,
cache: Rc<RefCell<Cache>>,
config: ExecutionManagerConfig,
inflight_checks: IndexMap<ClientOrderId, InflightCheck>,
external_order_claims: IndexMap<InstrumentId, StrategyId>,
processed_fills: RecencyMap<FillKey>,
recon_check_retries: IndexMap<ClientOrderId, u32>,
order_query_recency: RecencyMap<ClientOrderId>,
order_local_activity: RecencyMap<ClientOrderId>,
position_local_activity: RecencyMap<InstrumentAccountKey>,
position_local_activity_revisions: IndexMap<InstrumentAccountKey, u64>,
position_reconciliation_states: IndexMap<InstrumentAccountKey, PositionReconciliationState>,
position_reconciliation_tolerances: IndexMap<AccountId, Decimal>,
recent_fills_cache: RecencyMap<FillKey>,
missing_order_coverage_warnings: IndexSet<ClientOrderId>,
unresolved_order_coverage: IndexSet<ClientOrderId>,
targeted_order_queries: IndexSet<ClientOrderId>,
}
impl Debug for ExecutionManager {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct(stringify!(ExecutionManager))
.field("config", &self.config)
.field("inflight_checks", &self.inflight_checks)
.field("external_order_claims", &self.external_order_claims)
.field("processed_fills", &self.processed_fills)
.field("recon_check_retries", &self.recon_check_retries)
.finish_non_exhaustive()
}
}
impl ExecutionManager {
pub fn new(
clock: Rc<RefCell<dyn Clock>>,
cache: Rc<RefCell<Cache>>,
config: ExecutionManagerConfig,
) -> Self {
Self {
clock,
cache,
config,
inflight_checks: IndexMap::new(),
external_order_claims: IndexMap::new(),
processed_fills: RecencyMap::default(),
recon_check_retries: IndexMap::new(),
order_query_recency: RecencyMap::default(),
order_local_activity: RecencyMap::default(),
position_local_activity: RecencyMap::default(),
position_local_activity_revisions: IndexMap::new(),
position_reconciliation_states: IndexMap::new(),
position_reconciliation_tolerances: IndexMap::new(),
recent_fills_cache: RecencyMap::default(),
missing_order_coverage_warnings: IndexSet::new(),
unresolved_order_coverage: IndexSet::new(),
targeted_order_queries: IndexSet::new(),
}
}
pub(crate) fn set_position_reconciliation_tolerance(
&mut self,
account_id: AccountId,
tolerance: Decimal,
) {
let tolerance = if tolerance < Decimal::ZERO {
log::error!(
"Invalid negative position reconciliation tolerance {tolerance} for \
{account_id}; using the default"
);
DEFAULT_POSITION_RECONCILIATION_TOLERANCE
} else {
tolerance
};
self.position_reconciliation_tolerances
.insert(account_id, tolerance);
}
fn position_reconciliation_tolerance(&self, account_id: AccountId) -> Decimal {
self.position_reconciliation_tolerances
.get(&account_id)
.copied()
.unwrap_or(DEFAULT_POSITION_RECONCILIATION_TOLERANCE)
}
#[allow(unknown_lints, reason = "Clippy lint is unavailable on Rust 1.97")]
#[expect(
clippy::unused_async,
clippy::unused_async_trait_impl,
reason = "public reconciliation API stays async; live node and test callers await it"
)]
pub async fn reconcile_execution_mass_status(
&mut self,
mass_status: ExecutionMassStatus,
exec_engine: Rc<RefCell<ExecutionEngine>>,
) -> ReconciliationResult {
let raw_order_status_topic =
MessagingSwitchboard::reconciliation_raw_order_status_report_topic();
for report in mass_status.order_reports().values() {
msgbus::publish_any(raw_order_status_topic, report);
}
let raw_fill_topic = MessagingSwitchboard::reconciliation_raw_fill_report_topic();
for fills in mass_status.fill_reports().values() {
for fill in fills {
msgbus::publish_any(raw_fill_topic, fill);
}
}
let raw_position_topic =
MessagingSwitchboard::reconciliation_raw_position_status_report_topic();
for reports in mass_status.position_reports().values() {
for report in reports {
msgbus::publish_any(raw_position_topic, report);
}
}
let venue = mass_status.venue;
let order_count = mass_status.order_reports().len();
let fill_count: usize = mass_status.fill_reports().values().map(|v| v.len()).sum();
let position_count = mass_status.position_reports().len();
log_info!(
"Reconciling ExecutionMassStatus for {venue}",
color = LogColor::Blue
);
log_info!(
"Received {order_count} order(s), {fill_count} fill(s), {position_count} position(s)",
color = LogColor::Blue
);
let retained_fill_state = self.retained_fill_state();
let reported_fill_keys: IndexSet<(AccountId, InstrumentId, TradeId)> = mass_status
.fill_reports()
.values()
.flatten()
.map(|fill| (fill.account_id, fill.instrument_id, fill.trade_id))
.collect();
let (adjusted_order_reports, adjusted_fill_reports) =
self.adjust_mass_status_fills(&mass_status);
let mut events = Vec::new();
let mut external_orders = Vec::new();
let mut orders_reconciled = 0usize;
let mut external_orders_created = 0usize;
let mut open_orders_initialized = 0usize;
let mut orders_skipped_no_instrument = 0usize;
let mut orders_skipped_duplicate = 0usize;
let mut fills_applied = 0usize;
let mut fill_queue = ReconciliationFillQueue::default();
let fill_reports = &adjusted_fill_reports;
let mut seen_fill_keys: IndexSet<FillKey> = IndexSet::new();
for fills in fill_reports.values() {
for fill in fills {
let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
if !seen_fill_keys.insert(fill_key) {
log::warn!(
"Duplicate trade_id {} for {} in mass status",
fill.trade_id,
fill.instrument_id
);
}
}
}
let order_reports = Self::deduplicate_order_reports(adjusted_order_reports.values());
let mut orders_skipped_filtered = 0usize;
for report in order_reports.values() {
if self.should_skip_order_report(report) {
orders_skipped_filtered += 1;
continue;
}
if let Some(client_order_id) = &report.client_order_id {
if let Some(cached_order) = self.get_order(*client_order_id)
&& Self::is_exact_order_match(&cached_order, report)
{
log::debug!("Skipping order {client_order_id}: already in sync with venue");
orders_skipped_duplicate += 1;
if let Err(e) = self.cache.borrow_mut().add_venue_order_id(
client_order_id,
&report.venue_order_id,
false,
) {
log::warn!("Failed to add venue order ID index: {e}");
}
continue;
}
if let Some(cached_order) = self.get_order(*client_order_id)
&& cached_order.is_closed()
&& cached_order
.tags()
.is_some_and(|tags| tags.contains(&*TAG_RECONCILIATION))
{
log::debug!(
"Skipping closed reconciliation order {client_order_id}: \
synthetic position adjustment from previous session",
);
orders_skipped_duplicate += 1;
continue;
}
if let Some(order) = self.get_order(*client_order_id) {
let instrument = self.get_instrument(&report.instrument_id);
log::info!(
color = LogColor::Blue as u8;
"Reconciling {} {} {} [{}] -> [{}]",
client_order_id,
report.venue_order_id,
report.instrument_id,
order.status(),
report.order_status,
);
let order_fills: Vec<&FillReport> = fill_reports
.get(&report.venue_order_id)
.map(|f| f.iter().collect())
.unwrap_or_default();
let order_events = self.reconcile_order_with_fills(
&order,
report,
&order_fills,
instrument.as_ref(),
&mut fill_queue,
);
if !order_events.is_empty() {
orders_reconciled += 1;
fills_applied += order_events
.iter()
.filter(|e| matches!(e, OrderEventAny::Filled(_)))
.count();
events.extend(order_events);
}
if let Err(e) = self.cache.borrow_mut().add_venue_order_id(
client_order_id,
&report.venue_order_id,
false,
) {
log::warn!("Failed to add venue order ID index: {e}");
}
} else if let Some(order) = self.get_order_by_venue_order_id(report.venue_order_id)
{
let instrument = self.get_instrument(&report.instrument_id);
log::info!(
color = LogColor::Blue as u8;
"Reconciling {} (matched by venue_order_id {}) {} [{}] -> [{}]",
order.client_order_id(),
report.venue_order_id,
report.instrument_id,
order.status(),
report.order_status,
);
let order_fills: Vec<&FillReport> = fill_reports
.get(&report.venue_order_id)
.map(|f| f.iter().collect())
.unwrap_or_default();
let order_events = self.reconcile_order_with_fills(
&order,
report,
&order_fills,
instrument.as_ref(),
&mut fill_queue,
);
if !order_events.is_empty() {
orders_reconciled += 1;
fills_applied += order_events
.iter()
.filter(|e| matches!(e, OrderEventAny::Filled(_)))
.count();
events.extend(order_events);
}
if let Err(e) = self.cache.borrow_mut().add_venue_order_id(
&order.client_order_id(),
&report.venue_order_id,
false,
) {
log::warn!("Failed to add venue order ID index: {e}");
}
} else if !self.config.filter_unclaimed_external {
if let Some(instrument) = self.get_instrument(&report.instrument_id) {
let order_fills: Vec<&FillReport> = fill_reports
.get(&report.venue_order_id)
.map(|f| f.iter().collect())
.unwrap_or_default();
let (external_events, metadata) = self.handle_external_order(
report,
mass_status.account_id,
&instrument,
&order_fills,
false, Some(&mut fill_queue),
);
if !external_events.is_empty() {
external_orders_created += 1;
fills_applied += external_events
.iter()
.filter(|e| matches!(e, OrderEventAny::Filled(_)))
.count();
if report.order_status.is_open() {
open_orders_initialized += 1;
}
events.extend(external_events);
if let Some(m) = metadata {
external_orders.push(m);
}
}
} else {
orders_skipped_no_instrument += 1;
}
}
} else if let Some(order) = self.get_order_by_venue_order_id(report.venue_order_id) {
let instrument = self.get_instrument(&report.instrument_id);
log::info!(
color = LogColor::Blue as u8;
"Reconciling {} (matched by venue_order_id {}) {} [{}] -> [{}]",
order.client_order_id(),
report.venue_order_id,
report.instrument_id,
order.status(),
report.order_status,
);
let order_fills: Vec<&FillReport> = fill_reports
.get(&report.venue_order_id)
.map(|f| f.iter().collect())
.unwrap_or_default();
let order_events = self.reconcile_order_with_fills(
&order,
report,
&order_fills,
instrument.as_ref(),
&mut fill_queue,
);
if !order_events.is_empty() {
orders_reconciled += 1;
fills_applied += order_events
.iter()
.filter(|e| matches!(e, OrderEventAny::Filled(_)))
.count();
events.extend(order_events);
}
if let Err(e) = self.cache.borrow_mut().add_venue_order_id(
&order.client_order_id(),
&report.venue_order_id,
false,
) {
log::warn!("Failed to add venue order ID index: {e}");
}
} else if let Some(instrument) = self.get_instrument(&report.instrument_id) {
let is_synthetic = report.venue_order_id.as_str().starts_with("S-");
let order_fills: Vec<&FillReport> = fill_reports
.get(&report.venue_order_id)
.map(|f| f.iter().collect())
.unwrap_or_default();
let (external_events, metadata) = self.handle_external_order(
report,
mass_status.account_id,
&instrument,
&order_fills,
is_synthetic,
Some(&mut fill_queue),
);
if !external_events.is_empty() {
external_orders_created += 1;
fills_applied += external_events
.iter()
.filter(|e| matches!(e, OrderEventAny::Filled(_)))
.count();
if report.order_status.is_open() {
open_orders_initialized += 1;
}
events.extend(external_events);
if let Some(m) = metadata {
external_orders.push(m);
}
}
} else {
orders_skipped_no_instrument += 1;
}
}
let processed_venue_order_ids: IndexSet<VenueOrderId> =
order_reports.keys().copied().collect();
for (venue_order_id, fills) in fill_reports {
if processed_venue_order_ids.contains(venue_order_id) {
continue;
}
let Some(first_fill) = fills.first() else {
continue;
};
if !self.should_reconcile_instrument(&first_fill.instrument_id) {
log::debug!(
"Skipping orphan fills for {}: not in reconciliation_instrument_ids",
first_fill.instrument_id
);
continue;
}
if let Some(client_order_id) = &first_fill.client_order_id
&& self
.config
.filtered_client_order_ids
.contains(client_order_id)
{
log::debug!(
"Skipping orphan fills for {client_order_id}: in filtered_client_order_ids"
);
continue;
}
let order = first_fill
.client_order_id
.as_ref()
.and_then(|id| self.get_order(*id))
.or_else(|| self.get_order_by_venue_order_id(*venue_order_id));
if let Some(ref order) = order
&& self
.config
.filtered_client_order_ids
.contains(&order.client_order_id())
{
log::debug!(
"Skipping orphan fills for {}: in filtered_client_order_ids",
order.client_order_id()
);
continue;
}
if let Some(order) = order {
let instrument_id = order.instrument_id();
if let Some(instrument) = self.get_instrument(&instrument_id) {
let mut sorted_fills: Vec<&FillReport> = fills.iter().collect();
sorted_fills.sort_by_key(|f| f.ts_event);
for fill in sorted_fills {
if let Some((event, fill_key)) = self.create_order_fill(
&order,
fill,
&instrument,
&fill_queue.pending_fill_keys,
) {
fills_applied += 1;
fill_queue.push(&mut events, event, fill_key);
}
}
}
}
}
events.sort_by_key(|e| e.ts_event());
for event in &events {
if let OrderEventAny::Filled(fill) = event
&& (retained_fill_state.fill_keys.contains(&(
fill.account_id,
fill.instrument_id,
fill.trade_id,
)) || ((retained_fill_state.missing_order_ids.contains(&(
fill.account_id,
fill.instrument_id,
fill.client_order_id,
)) || retained_fill_state.missing_venue_order_ids.contains(&(
fill.account_id,
fill.instrument_id,
fill.venue_order_id,
))) && !reported_fill_keys.contains(&(
fill.account_id,
fill.instrument_id,
fill.trade_id,
))) || retained_fill_state
.netting_lifecycle_starts
.get(&(fill.account_id, fill.instrument_id, fill.strategy_id))
.is_some_and(|ts_opened| fill.ts_event < *ts_opened))
{
exec_engine.borrow_mut().project_reconciliation_fill(fill);
} else {
exec_engine.borrow_mut().process(event);
}
if let OrderEventAny::Filled(fill) = event
&& let Some(fill_key) = fill_queue.event_fill_keys.get(&fill.event_id).copied()
&& self.is_fill_applied(fill, fill_key)
{
self.processed_fills.mark(fill_key);
}
}
let mut positions_created = 0usize;
if !self.config.filter_position_reports {
let instruments_with_unattributed_fills: IndexSet<InstrumentId> = mass_status
.fill_reports()
.values()
.flatten()
.filter(|f| f.venue_position_id.is_none())
.map(|f| f.instrument_id)
.chain(
mass_status
.order_reports()
.values()
.filter(|r| !r.filled_qty.is_zero() && r.venue_position_id.is_none())
.map(|r| r.instrument_id),
)
.collect();
let positions_with_fills: IndexSet<PositionId> = mass_status
.fill_reports()
.values()
.flatten()
.filter_map(|f| f.venue_position_id)
.chain(
mass_status
.order_reports()
.values()
.filter(|r| !r.filled_qty.is_zero())
.filter_map(|r| r.venue_position_id),
)
.collect();
for (instrument_id, reports) in mass_status.position_reports() {
if !self.should_reconcile_instrument(&instrument_id) {
log::debug!(
"Skipping position reports for {instrument_id}: not in reconciliation_instrument_ids"
);
continue;
}
for report in reports {
if let Some(position_events) = self.reconcile_position_report(
&report,
mass_status.account_id,
&instruments_with_unattributed_fills,
&positions_with_fills,
) {
for event in position_events {
exec_engine.borrow_mut().process(&event);
events.push(event);
}
positions_created += 1;
}
}
}
}
if orders_skipped_no_instrument > 0 {
log::warn!("{orders_skipped_no_instrument} orders skipped (instrument not in cache)");
}
if orders_skipped_duplicate > 0 {
log::debug!("{orders_skipped_duplicate} orders skipped (already in sync)");
}
if orders_skipped_filtered > 0 {
log::debug!("{orders_skipped_filtered} orders skipped (filtered by config)");
}
log::info!(
color = LogColor::Blue as u8;
"Reconciliation complete for {venue}: reconciled={orders_reconciled}, external={external_orders_created}, open={open_orders_initialized}, fills={fills_applied}, positions={positions_created}, skipped={orders_skipped_duplicate}, filtered={orders_skipped_filtered}",
);
ReconciliationResult {
events,
external_orders,
}
}
fn retained_fill_state(&self) -> RetainedFillState {
let cache = self.cache.borrow();
let positions = cache.positions(None, None, None, None, None);
let mut fill_keys = IndexSet::new();
let mut missing_order_ids = IndexSet::new();
let mut missing_venue_order_ids = IndexSet::new();
let mut netting_lifecycle_starts = IndexMap::new();
for position in positions {
for fill in &position.events {
fill_keys.insert((position.account_id, position.instrument_id, fill.trade_id));
if cache.order(&fill.client_order_id).is_none() {
missing_order_ids.insert((
position.account_id,
position.instrument_id,
fill.client_order_id,
));
missing_venue_order_ids.insert((
position.account_id,
position.instrument_id,
fill.venue_order_id,
));
}
}
if cache.oms_type(&position.id) == Some(OmsType::Netting) {
netting_lifecycle_starts.insert(
(
position.account_id,
position.instrument_id,
position.strategy_id,
),
position.ts_opened,
);
}
}
RetainedFillState {
fill_keys,
missing_order_ids,
missing_venue_order_ids,
netting_lifecycle_starts,
}
}
pub fn check_inflight_orders(&mut self) -> InflightCheckResult {
let mut result = InflightCheckResult::default();
let now = dst::time::Instant::now();
let threshold = Duration::from_millis(self.config.inflight_threshold_ms);
let mut to_check = Vec::new();
for (client_order_id, check) in &self.inflight_checks {
if now
.checked_duration_since(check.submitted_at)
.is_some_and(|elapsed| elapsed > threshold)
{
to_check.push(*client_order_id);
}
}
for client_order_id in to_check {
if self
.config
.filtered_client_order_ids
.contains(&client_order_id)
{
self.clear_recon_tracking(&client_order_id, true);
continue;
}
if self.targeted_order_queries.contains(&client_order_id) {
continue;
}
if let Some(check) = self.inflight_checks.get_mut(&client_order_id) {
if let Some(last_query_at) = check.last_query_at
&& now
.checked_duration_since(last_query_at)
.is_none_or(|elapsed| elapsed < threshold)
{
continue;
}
check.retry_count += 1;
check.last_query_at = Some(now);
self.order_query_recency.mark(client_order_id);
self.recon_check_retries
.insert(client_order_id, check.retry_count);
if check.retry_count >= self.config.inflight_max_retries {
let ts_now = self.clock.borrow().timestamp_ns();
if let Some(order) = self.get_order(client_order_id) {
match order.status() {
OrderStatus::Submitted => {
if let Some(event) = create_reconciliation_rejected(
&order,
Some("INFLIGHT_TIMEOUT"),
ts_now,
) {
result.events.push(event);
}
}
OrderStatus::PendingUpdate | OrderStatus::PendingCancel => {
let event = OrderEventAny::Canceled(OrderCanceled::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
UUID4::new(),
ts_now,
ts_now,
true, order.venue_order_id(),
order.account_id(),
));
result.events.push(event);
}
_ => {
}
}
}
self.clear_recon_tracking(&client_order_id, true);
} else if let Some(order) = self.get_order(client_order_id) {
let ts_now = self.clock.borrow().timestamp_ns();
let client_id = self.cache.borrow().client_id(&client_order_id).copied();
let query = TradingCommand::QueryOrder(QueryOrder::new(
order.trader_id(),
client_id,
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
order.venue_order_id(),
UUID4::new(),
ts_now,
None,
None, ));
result.queries.push(query);
}
}
}
result
}
fn filtered_open_orders_for_reconciliation(&self) -> Vec<OrderAny> {
{
let cache = self.cache.borrow();
let mut orders = cache.orders_open(None, None, None, None, None);
orders.extend(cache.orders_inflight(None, None, None, None, None));
let mut seen_client_order_ids = IndexSet::new();
orders.retain(|order| seen_client_order_ids.insert(order.client_order_id()));
if self.config.reconciliation_instrument_ids.is_empty() {
orders.iter().map(|o| (*o).clone()).collect()
} else {
orders
.iter()
.filter(|o| {
self.config
.reconciliation_instrument_ids
.contains(&o.instrument_id())
})
.map(|o| (*o).clone())
.collect()
}
}
}
fn open_position_keys_for_reconciliation(&self) -> IndexSet<InstrumentAccountKey> {
let cache = self.cache.borrow();
let positions = cache.positions_open(None, None, None, None, None);
let mut position_keys = IndexSet::new();
for position in positions {
if !self.should_reconcile_instrument(&position.instrument_id) {
continue;
}
position_keys.insert((position.instrument_id, position.account_id));
}
position_keys
}
pub(crate) fn prepare_open_order_report_check(
&mut self,
command_id: UUID4,
clients: &[&dyn ExecutionClient],
) -> OpenOrderReportCheck {
let filtered_orders = self.filtered_open_orders_for_reconciliation();
let active_order_ids: IndexSet<ClientOrderId> = filtered_orders
.iter()
.map(|order| order.client_order_id())
.collect();
self.missing_order_coverage_warnings
.retain(|client_order_id| active_order_ids.contains(client_order_id));
self.unresolved_order_coverage
.retain(|client_order_id| active_order_ids.contains(client_order_id));
let mut client_coverage = IndexMap::new();
for order in &filtered_orders {
let client_order_id = order.client_order_id();
let coverage = self.resolve_order_report_client_coverage(order, clients);
match &coverage {
ReportClientCoverage::Resolved(_) => {
if self
.unresolved_order_coverage
.shift_remove(&client_order_id)
{
self.missing_order_coverage_warnings
.shift_remove(&client_order_id);
}
}
ReportClientCoverage::Unresolved => {
self.unresolved_order_coverage.insert(client_order_id);
}
}
client_coverage.insert(client_order_id, coverage);
}
log::debug!(
"Found {} order{} open in cache",
filtered_orders.len(),
if filtered_orders.len() == 1 { "" } else { "s" }
);
let ts_now = self.clock.borrow().timestamp_ns();
let start = self.config.open_check_lookback_mins.map(|mins| {
let lookback_ns = mins_to_nanos(mins);
ts_now.saturating_sub_ns(lookback_ns)
});
let mut command = GenerateOrderStatusReports::new(
command_id,
ts_now,
self.config.open_check_open_only,
None,
start,
None,
None,
None,
);
command.log_receipt_level = LogLevel::Debug;
OpenOrderReportCheck {
command,
filtered_orders,
client_coverage,
start,
}
}
fn resolve_order_report_client_coverage(
&self,
order: &OrderAny,
clients: &[&dyn ExecutionClient],
) -> ReportClientCoverage {
if let Some(client_id) = self.cache.borrow().client_id(&order.client_order_id()) {
return ReportClientCoverage::Resolved(IndexSet::from([*client_id]));
}
if let Some(account_id) = order.account_id() {
let account_clients = clients
.iter()
.filter(|client| client.account_id() == account_id)
.map(|client| client.client_id())
.collect::<IndexSet<_>>();
if !account_clients.is_empty() {
return ReportClientCoverage::Resolved(account_clients);
}
}
let venue_clients = clients
.iter()
.filter(|client| client.handles_order_venue(order.instrument_id().venue))
.map(|client| client.client_id())
.collect::<IndexSet<_>>();
if venue_clients.is_empty() {
ReportClientCoverage::Unresolved
} else {
ReportClientCoverage::Resolved(venue_clients)
}
}
pub fn check_open_order_queries(&mut self) -> Vec<TradingCommand> {
self.check_open_order_queries_for_clients(None)
}
pub(crate) fn check_open_order_queries_for_clients(
&mut self,
client_ids: Option<&IndexSet<ClientId>>,
) -> Vec<TradingCommand> {
let now = dst::time::Instant::now();
let query_delay = Duration::from_millis(u64::from(self.config.single_order_query_delay_ms));
let query_limit = self.config.max_single_order_queries_per_cycle as usize;
if query_limit == 0 {
return Vec::new();
}
let mut filtered_orders = self.filtered_open_orders_for_reconciliation();
filtered_orders.sort_by_key(|order| {
let client_order_id = order.client_order_id();
(
self.order_query_recency.last_marked(&client_order_id),
client_order_id,
)
});
let mut queries = Vec::new();
for order in filtered_orders {
if queries.len() >= query_limit {
break;
}
let client_order_id = order.client_order_id();
let client_id = self.cache.borrow().client_id(&client_order_id).copied();
if let Some(client_ids) = client_ids
&& !client_id.is_some_and(|client_id| client_ids.contains(&client_id))
{
continue;
}
if self
.config
.filtered_client_order_ids
.contains(&client_order_id)
{
continue;
}
let threshold = Duration::from_nanos(self.config.open_check_threshold_ns);
if let Some(elapsed) = self.order_local_activity.elapsed_at(&client_order_id, now)
&& elapsed < threshold
{
let elapsed_ms = elapsed.as_millis();
let threshold_ms = threshold.as_millis();
log::debug!(
"Deferring open order query for {client_order_id}: recent local activity \
({elapsed_ms}ms < threshold={threshold_ms}ms)",
);
continue;
}
if self
.order_query_recency
.within_at(&client_order_id, now, query_delay)
{
continue;
}
self.order_query_recency.mark(client_order_id);
let ts_now = self.clock.borrow().timestamp_ns();
let cmd = TradingCommand::QueryOrder(QueryOrder::new(
order.trader_id(),
client_id,
order.strategy_id(),
order.instrument_id(),
client_order_id,
order.venue_order_id(),
UUID4::new(),
ts_now,
None,
None,
));
queries.push(cmd);
}
queries
}
pub async fn check_open_orders(
&mut self,
clients: &[&dyn ExecutionClient],
) -> Vec<OrderEventAny> {
log::debug!("Checking order consistency between cached-state and venues");
let check = self.prepare_open_order_report_check(UUID4::new(), clients);
let mut all_reports = Vec::new();
let mut queried_clients = IndexSet::new();
let mut failed_clients = IndexSet::new();
for client in clients {
let client_id = client.client_id();
queried_clients.insert(client_id);
match client.generate_order_status_reports(&check.command).await {
Ok(reports) => {
all_reports.extend(reports);
}
Err(e) => {
failed_clients.insert(client_id);
log::warn!(
"Failed to query order reports from {}: {e}",
client.client_id()
);
}
}
}
let result = self.reconcile_open_order_reports(
&check,
all_reports,
&queried_clients,
&failed_clients,
);
let mut events = result.events;
if !result.targeted_queries.is_empty() {
let query_delay =
Duration::from_millis(u64::from(self.config.single_order_query_delay_ms));
let query_results =
request_targeted_order_reports(clients, result.targeted_queries, query_delay).await;
events.extend(self.reconcile_targeted_order_reports(query_results));
}
events
}
pub(crate) fn reconcile_open_order_reports(
&mut self,
check: &OpenOrderReportCheck,
all_reports: Vec<OrderStatusReport>,
queried_clients: &IndexSet<ClientId>,
failed_clients: &IndexSet<ClientId>,
) -> OpenOrderReconciliationResult {
let mut venue_reported_ids = IndexSet::new();
for report in &all_reports {
if let Some(client_order_id) = &report.client_order_id {
venue_reported_ids.insert(*client_order_id);
self.missing_order_coverage_warnings
.shift_remove(client_order_id);
self.recon_check_retries.shift_remove(client_order_id);
} else {
let mapped_client_order_id = self
.cache
.borrow()
.client_order_id(&report.venue_order_id)
.copied();
if let Some(client_order_id) = mapped_client_order_id {
venue_reported_ids.insert(client_order_id);
self.missing_order_coverage_warnings
.shift_remove(&client_order_id);
self.recon_check_retries.shift_remove(&client_order_id);
}
}
}
let mut events = Vec::new();
let mut targeted_candidates = Vec::new();
for report in all_reports {
if let Some(client_order_id) = &report.client_order_id
&& let Some(order) = self.get_order(*client_order_id)
{
let threshold = Duration::from_nanos(self.config.open_check_threshold_ns);
if let Some(elapsed) = self.order_local_activity.elapsed(client_order_id)
&& elapsed < threshold
{
let elapsed_ms = elapsed.as_millis();
let threshold_ms = threshold.as_millis();
log::debug!(
"Deferring reconciliation for {client_order_id}: recent local activity ({elapsed_ms}ms < threshold={threshold_ms}ms)",
);
continue;
}
let instrument = self.get_instrument(&report.instrument_id);
if let Some(event) =
self.reconcile_order_report(&order, &report, instrument.as_ref())
{
events.push(event);
}
}
}
if self.config.open_check_open_only {
let cached_ids: IndexSet<ClientOrderId> = check
.filtered_orders
.iter()
.map(|o| o.client_order_id())
.collect();
let missing_at_venue: IndexSet<ClientOrderId> = cached_ids
.difference(&venue_reported_ids)
.copied()
.collect();
if !missing_at_venue.is_empty() {
log::debug!(
"{} cached open order{} not present in venue current response",
missing_at_venue.len(),
if missing_at_venue.len() == 1 {
" is"
} else {
"s are"
},
);
for client_order_id in missing_at_venue {
log::debug!("Cached open order missing from venue response: {client_order_id}");
}
}
} else {
let candidates: Vec<&OrderAny> = if let Some(cutoff) = check.start {
check
.filtered_orders
.iter()
.filter(|o| o.ts_last() >= cutoff)
.collect()
} else {
check.filtered_orders.iter().collect()
};
for order in candidates {
let client_order_id = order.client_order_id();
if venue_reported_ids.contains(&client_order_id) {
continue;
}
let coverage = check
.client_coverage
.get(&client_order_id)
.unwrap_or(&ReportClientCoverage::Unresolved);
let ReportClientCoverage::Resolved(responsible_clients) = coverage else {
if self.missing_order_coverage_warnings.insert(client_order_id) {
log::warn!(
"Skipping order reconciliation for {client_order_id}: responsible execution client coverage is unresolved"
);
}
continue;
};
if responsible_clients.is_empty() {
if self.missing_order_coverage_warnings.insert(client_order_id) {
log::warn!(
"Skipping order reconciliation for {client_order_id}: responsible execution client coverage is unresolved"
);
}
continue;
}
let missing_clients = responsible_clients
.difference(queried_clients)
.copied()
.collect::<IndexSet<_>>();
if !missing_clients.is_empty() {
if self.missing_order_coverage_warnings.insert(client_order_id) {
log::warn!(
"Skipping order reconciliation for {client_order_id}: responsible execution clients were not queried: {missing_clients:?}"
);
}
continue;
}
let failed_responsible_clients = responsible_clients
.intersection(failed_clients)
.copied()
.collect::<IndexSet<_>>();
if !failed_responsible_clients.is_empty() {
log::warn!(
"Skipping order reconciliation for {client_order_id}: failed to query responsible execution clients: {failed_responsible_clients:?}"
);
continue;
}
self.missing_order_coverage_warnings
.shift_remove(&client_order_id);
if let Some(order) = self.prepare_missing_order_query(client_order_id) {
targeted_candidates.push((order, responsible_clients.clone()));
}
}
}
targeted_candidates.sort_by_key(|(order, _)| {
let client_order_id = order.client_order_id();
(
self.order_query_recency.last_marked(&client_order_id),
client_order_id,
)
});
let query_limit = self.config.max_single_order_queries_per_cycle as usize;
let mut planned_queries = 0usize;
let mut cap_deferred_orders = 0usize;
let mut targeted_queries = Vec::new();
for (order, responsible_clients) in targeted_candidates {
let client_order_id = order.client_order_id();
let required_queries = responsible_clients.len();
let exceeds_query_limit = planned_queries + required_queries > query_limit;
let can_run_oversized_group = planned_queries == 0 && query_limit > 0;
if required_queries == 0 || (exceeds_query_limit && !can_run_oversized_group) {
cap_deferred_orders += 1;
continue;
}
if required_queries > query_limit {
log::warn!(
"Targeted order query for {client_order_id} requires {required_queries} responsible clients, exceeding the per-cycle limit {query_limit} to avoid indefinite deferral"
);
}
planned_queries += required_queries;
self.order_query_recency.mark(client_order_id);
self.targeted_order_queries.insert(client_order_id);
targeted_queries.push(TargetedOrderQuery {
client_order_id,
responsible_clients,
command: GenerateOrderStatusReport::new(
UUID4::new(),
self.clock.borrow().timestamp_ns(),
Some(order.instrument_id()),
Some(client_order_id),
order.venue_order_id(),
None,
None,
),
});
}
if cap_deferred_orders > 0 {
log::warn!(
"Reached max single-order queries ({query_limit}) this cycle, deferring {cap_deferred_orders} order(s)"
);
}
OpenOrderReconciliationResult {
events,
targeted_queries,
}
}
pub(crate) fn reconcile_targeted_order_reports(
&mut self,
results: Vec<TargetedOrderReportResult>,
) -> Vec<OrderEventAny> {
let mut events = Vec::new();
for result in results {
let client_order_id = result.client_order_id;
self.targeted_order_queries.shift_remove(&client_order_id);
if let Some(report) = result.report {
self.recon_check_retries.shift_remove(&client_order_id);
self.missing_order_coverage_warnings
.shift_remove(&client_order_id);
let Some(order) = self.get_order(client_order_id) else {
continue;
};
let instrument = self.get_instrument(&report.instrument_id);
log::info!(
color = LogColor::Blue as u8;
"Found {client_order_id} via targeted order status query: {}",
report.order_status,
);
if let Some(event) =
self.reconcile_order_report(&order, &report, instrument.as_ref())
{
events.push(event);
}
continue;
}
if result.coverage_complete {
events.extend(self.resolve_missing_order(client_order_id));
} else {
log::warn!(
"Deferring missing-order resolution for {client_order_id}: targeted order status coverage was incomplete"
);
}
}
events
}
#[must_use]
pub(crate) fn prepare_position_report_check(
&self,
command_id: UUID4,
clients: &[&dyn ExecutionClient],
) -> PositionReportCheck {
let position_keys = self.open_position_keys_for_reconciliation();
let client_coverage = position_keys
.iter()
.map(|key| {
(
*key,
Self::resolve_position_report_client_coverage(*key, clients),
)
})
.collect();
let activity_revisions = position_keys
.iter()
.map(|key| (*key, self.position_activity_revision(key)))
.collect();
log::debug!(
"Found {} unique instrument/account combination{} with open positions",
position_keys.len(),
if position_keys.len() == 1 { "" } else { "s" }
);
let mut command = GeneratePositionStatusReports::new(
command_id,
self.clock.borrow().timestamp_ns(),
None, None, None, None, None, );
command.log_receipt_level = LogLevel::Debug;
PositionReportCheck {
command,
client_coverage,
activity_revisions,
}
}
fn resolve_position_report_client_coverage(
key: InstrumentAccountKey,
clients: &[&dyn ExecutionClient],
) -> ReportClientCoverage {
let account_clients = clients
.iter()
.filter(|client| client.account_id() == key.1)
.map(|client| client.client_id())
.collect::<IndexSet<_>>();
if !account_clients.is_empty() {
return ReportClientCoverage::Resolved(account_clients);
}
let venue_clients = clients
.iter()
.filter(|client| client.handles_order_venue(key.0.venue))
.map(|client| client.client_id())
.collect::<IndexSet<_>>();
if venue_clients.is_empty() {
ReportClientCoverage::Unresolved
} else {
ReportClientCoverage::Resolved(venue_clients)
}
}
pub async fn check_positions_consistency(
&mut self,
clients: &[&dyn ExecutionClient],
) -> Vec<OrderEventAny> {
let check = self.prepare_position_report_check(UUID4::new(), clients);
let mut reports = Vec::new();
let mut queried_clients = IndexSet::new();
let mut failed_clients = IndexSet::new();
for client in clients {
let client_id = client.client_id();
queried_clients.insert(client_id);
self.set_position_reconciliation_tolerance(
client.account_id(),
client.position_reconciliation_tolerance(),
);
match client
.generate_position_status_reports(&check.command)
.await
{
Ok(client_reports) => {
reports.extend(client_reports);
}
Err(e) => {
failed_clients.insert(client_id);
log::warn!(
"Failed to query position reports from {}: {e}",
client.client_id()
);
}
}
}
self.reconcile_position_reports(&check, reports, &queried_clients, &failed_clients)
}
#[must_use]
pub(crate) fn reconcile_position_reports(
&mut self,
check: &PositionReportCheck,
reports: Vec<PositionStatusReport>,
queried_clients: &IndexSet<ClientId>,
failed_clients: &IndexSet<ClientId>,
) -> Vec<OrderEventAny> {
log::debug!("Checking position consistency between cached-state and venues");
let mut venue_positions: IndexMap<InstrumentAccountKey, Vec<PositionStatusReport>> =
IndexMap::new();
for report in reports {
if !self.should_reconcile_instrument(&report.instrument_id) {
continue;
}
venue_positions
.entry((report.instrument_id, report.account_id))
.or_default()
.push(report);
}
let mut events = Vec::new();
for key in check.client_coverage.keys() {
let prepared_revision = check
.activity_revisions
.get(key)
.copied()
.unwrap_or_default();
if self.position_activity_revision(key) > prepared_revision {
log::debug!(
"Deferring position reconciliation for {}/{}: local activity recorded during report request",
key.0,
key.1,
);
continue;
}
let venue_reports = venue_positions
.get(key)
.map(Vec::as_slice)
.unwrap_or_default();
if venue_reports.is_empty() {
match check.client_coverage.get(key) {
Some(ReportClientCoverage::Resolved(responsible_clients))
if !responsible_clients.is_empty()
&& responsible_clients.is_subset(queried_clients)
&& responsible_clients.is_disjoint(failed_clients) => {}
Some(ReportClientCoverage::Resolved(responsible_clients))
if responsible_clients.is_empty() =>
{
log::warn!(
"Skipping position reconciliation for {}/{}: responsible execution client coverage is unresolved",
key.0,
key.1,
);
continue;
}
Some(ReportClientCoverage::Resolved(responsible_clients))
if !responsible_clients.is_subset(queried_clients) =>
{
log::warn!(
"Skipping position reconciliation for {}/{}: responsible execution clients were not all queried",
key.0,
key.1,
);
continue;
}
Some(ReportClientCoverage::Resolved(responsible_clients)) => {
let failed_responsible_clients = responsible_clients
.intersection(failed_clients)
.copied()
.collect::<IndexSet<_>>();
log::warn!(
"Skipping position reconciliation for {}/{}: failed to query responsible execution clients: {failed_responsible_clients:?}",
key.0,
key.1,
);
continue;
}
Some(ReportClientCoverage::Unresolved) | None => {
log::warn!(
"Skipping position reconciliation for {}/{}: responsible execution client coverage is unresolved",
key.0,
key.1,
);
continue;
}
}
}
if let Some(discrepancy_events) = self.check_position_discrepancy(*key, venue_reports) {
events.extend(discrepancy_events);
}
}
let current_position_keys = self.open_position_keys_for_reconciliation();
for (key, venue_reports) in &venue_positions {
if check.client_coverage.contains_key(key)
|| venue_reports
.iter()
.all(|report| report.signed_decimal_qty == Decimal::ZERO)
{
continue;
}
if current_position_keys.contains(key) {
log::debug!(
"Deferring position reconciliation for {}/{}: position opened after client coverage was recorded",
key.0,
key.1,
);
continue;
}
if let Some(discrepancy_events) = self.check_position_discrepancy(*key, venue_reports) {
events.extend(discrepancy_events);
}
}
let active_keys: IndexSet<InstrumentAccountKey> = current_position_keys
.into_iter()
.chain(
venue_positions
.iter()
.filter(|(_, reports)| {
reports
.iter()
.any(|report| report.signed_decimal_qty != Decimal::ZERO)
})
.map(|(k, _)| *k),
)
.collect();
self.position_reconciliation_states
.retain(|k, _| active_keys.contains(k));
events
}
fn positions_avg_px(cached_positions: &[Position]) -> Option<Decimal> {
let mut total_value = Decimal::ZERO;
let mut total_qty = Decimal::ZERO;
for position in cached_positions {
let qty = position.signed_decimal_qty().abs();
if position.avg_px_open > 0.0
&& qty > Decimal::ZERO
&& let Ok(avg_px) = Decimal::from_str(&position.avg_px_open.to_string())
{
total_value += avg_px * qty;
total_qty += qty;
}
}
if total_qty > Decimal::ZERO {
Some(total_value / total_qty)
} else {
None
}
}
pub fn register_inflight(&mut self, client_order_id: ClientOrderId) {
if self
.config
.filtered_client_order_ids
.contains(&client_order_id)
{
return;
}
self.inflight_checks.insert(
client_order_id,
InflightCheck {
client_order_id,
submitted_at: dst::time::Instant::now(),
retry_count: 0,
last_query_at: None,
},
);
self.recon_check_retries.insert(client_order_id, 0);
self.order_query_recency.remove(&client_order_id);
self.order_local_activity.remove(&client_order_id);
}
pub fn record_local_activity(&mut self, client_order_id: ClientOrderId) {
self.order_local_activity.mark(client_order_id);
}
pub fn clear_recon_tracking(&mut self, client_order_id: &ClientOrderId, drop_last_query: bool) {
self.inflight_checks.shift_remove(client_order_id);
self.recon_check_retries.shift_remove(client_order_id);
self.missing_order_coverage_warnings
.shift_remove(client_order_id);
self.unresolved_order_coverage.shift_remove(client_order_id);
self.targeted_order_queries.shift_remove(client_order_id);
if drop_last_query {
self.order_query_recency.remove(client_order_id);
}
self.order_local_activity.remove(client_order_id);
}
#[cfg(feature = "node")]
pub(crate) fn remove_targeted_order_queries(&mut self, client_order_ids: &[ClientOrderId]) {
for client_order_id in client_order_ids {
self.targeted_order_queries.shift_remove(client_order_id);
}
}
#[must_use]
pub fn get_external_order_claim(&self, instrument_id: &InstrumentId) -> Option<StrategyId> {
self.external_order_claims.get(instrument_id).copied()
}
pub fn claim_external_orders(
&mut self,
instrument_id: InstrumentId,
strategy_id: StrategyId,
) -> anyhow::Result<()> {
if let Some(existing) = self.external_order_claims.get(&instrument_id) {
anyhow::bail!("External order claim for {instrument_id} already exists for {existing}");
}
self.external_order_claims
.insert(instrument_id, strategy_id);
Ok(())
}
pub fn record_position_activity(&mut self, instrument_id: InstrumentId, account_id: AccountId) {
let key = (instrument_id, account_id);
self.position_local_activity.mark(key);
let revision = self
.position_local_activity_revisions
.entry(key)
.or_default();
*revision = revision.saturating_add(1);
}
fn position_activity_revision(&self, key: &InstrumentAccountKey) -> u64 {
self.position_local_activity_revisions
.get(key)
.copied()
.unwrap_or_default()
}
#[must_use]
pub fn position_recon_retry_count(&self, key: &InstrumentAccountKey) -> u32 {
self.position_reconciliation_states
.get(key)
.map_or(0, |state| state.retries)
}
#[must_use]
pub fn recon_check_retry_count(&self, client_order_id: &ClientOrderId) -> u32 {
self.recon_check_retries
.get(client_order_id)
.copied()
.unwrap_or(0)
}
pub fn observe_order_event(&mut self, event: &OrderEventAny) {
match event {
OrderEventAny::Filled(fill) => {
self.record_position_activity(fill.instrument_id, fill.account_id);
}
OrderEventAny::Accepted(_)
| OrderEventAny::Rejected(_)
| OrderEventAny::Canceled(_)
| OrderEventAny::Expired(_)
| OrderEventAny::Denied(_)
| OrderEventAny::Updated(_)
| OrderEventAny::ModifyRejected(_)
| OrderEventAny::CancelRejected(_) => {
self.clear_recon_tracking(&event.client_order_id(), true);
}
_ => {}
}
self.record_local_activity(event.client_order_id());
}
pub fn observe_execution_report(&mut self, report: &ExecutionReport) {
match report {
ExecutionReport::Order(order_report) => {
self.observe_order_status_report(order_report);
}
ExecutionReport::Fill(fill_report) => {
let client_order_id = fill_report.client_order_id.or_else(|| {
self.cache
.borrow()
.client_order_id(&fill_report.venue_order_id)
.copied()
});
if let Some(coid) = client_order_id {
self.record_local_activity(coid);
}
self.record_position_activity(fill_report.instrument_id, fill_report.account_id);
}
ExecutionReport::OrderWithFills(order_report, fills) => {
self.observe_order_status_report(order_report);
for fill_report in fills {
self.record_position_activity(
fill_report.instrument_id,
fill_report.account_id,
);
}
}
ExecutionReport::Position(position_report) => {
self.record_position_activity(
position_report.instrument_id,
position_report.account_id,
);
}
ExecutionReport::MassStatus(_) => {
}
}
}
fn observe_order_status_report(&mut self, report: &OrderStatusReport) {
let Some(client_order_id) = report.client_order_id else {
return;
};
if !matches!(
report.order_status,
OrderStatus::PendingUpdate | OrderStatus::PendingCancel
) {
self.clear_recon_tracking(&client_order_id, report.order_status.is_closed());
}
self.record_local_activity(client_order_id);
}
#[must_use]
pub fn is_fill_recently_processed(
&self,
account_id: AccountId,
instrument_id: InstrumentId,
trade_id: TradeId,
) -> bool {
self.recent_fills_cache
.contains_key(&(account_id, instrument_id, trade_id))
}
pub fn mark_fill_processed(
&mut self,
account_id: AccountId,
instrument_id: InstrumentId,
trade_id: TradeId,
) {
self.recent_fills_cache
.mark((account_id, instrument_id, trade_id));
}
pub fn commit_recent_fill_if_applied(&mut self, fill: &OrderFilled) {
let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
if self.is_fill_applied(fill, fill_key) {
self.mark_fill_processed(fill_key.0, fill_key.1, fill_key.2);
}
}
pub fn prune_recent_fills_cache(&mut self, ttl_secs: f64) {
let ttl = match Duration::try_from_secs_f64(ttl_secs) {
Ok(ttl) => ttl,
Err(_) if ttl_secs > 0.0 => Duration::MAX,
Err(_) => Duration::ZERO,
};
self.recent_fills_cache.prune_older_than(ttl);
}
pub fn prune_processed_fills(&mut self) {
let Some(lookback_mins) = self.config.lookback_mins else {
return;
};
let ttl = Duration::from_mins(lookback_mins).max(Duration::from_mins(1));
self.processed_fills.prune_older_than(ttl);
}
pub fn prune_order_local_activity(&mut self) {
self.order_local_activity
.prune_older_than(Duration::from_nanos(self.config.open_check_threshold_ns));
}
pub fn purge_closed_orders(&mut self) {
let Some(buffer_mins) = self.config.purge_closed_orders_buffer_mins else {
return;
};
let ts_now = self.clock.borrow().timestamp_ns();
let buffer_secs = mins_to_secs(u64::from(buffer_mins));
self.cache
.borrow_mut()
.purge_closed_orders(ts_now, buffer_secs);
}
pub fn purge_closed_positions(&mut self) {
let Some(buffer_mins) = self.config.purge_closed_positions_buffer_mins else {
return;
};
let ts_now = self.clock.borrow().timestamp_ns();
let buffer_secs = mins_to_secs(u64::from(buffer_mins));
self.cache
.borrow_mut()
.purge_closed_positions(ts_now, buffer_secs);
}
pub fn purge_account_events(&mut self) {
let Some(lookback_mins) = self.config.purge_account_events_lookback_mins else {
return;
};
let ts_now = self.clock.borrow().timestamp_ns();
let lookback_secs = mins_to_secs(u64::from(lookback_mins));
self.cache
.borrow_mut()
.purge_account_events(ts_now, lookback_secs);
}
fn get_order(&self, client_order_id: ClientOrderId) -> Option<OrderAny> {
self.cache
.borrow()
.order(&client_order_id)
.map(|o| o.clone())
}
fn get_order_by_venue_order_id(&self, venue_order_id: VenueOrderId) -> Option<OrderAny> {
let cache = self.cache.borrow();
cache
.client_order_id(&venue_order_id)
.and_then(|client_order_id| cache.order(client_order_id).map(|o| o.clone()))
}
fn get_instrument(&self, instrument_id: &InstrumentId) -> Option<InstrumentAny> {
self.cache.borrow().instrument(instrument_id).cloned()
}
fn should_skip_order_report(&self, report: &OrderStatusReport) -> bool {
if let Some(client_order_id) = &report.client_order_id
&& self
.config
.filtered_client_order_ids
.contains(client_order_id)
{
log::debug!(
"Skipping order report {client_order_id}: in filtered_client_order_ids list"
);
return true;
}
if !self.should_reconcile_instrument(&report.instrument_id) {
log::debug!(
"Skipping order report for {}: not in reconciliation_instrument_ids",
report.instrument_id
);
return true;
}
false
}
fn should_reconcile_instrument(&self, instrument_id: &InstrumentId) -> bool {
self.config.reconciliation_instrument_ids.is_empty()
|| self
.config
.reconciliation_instrument_ids
.contains(instrument_id)
}
fn prepare_missing_order_query(&mut self, client_order_id: ClientOrderId) -> Option<OrderAny> {
let order = self.get_order(client_order_id)?;
if order.status().is_closed() {
log::debug!(
"Skipping missing-order resolution for {client_order_id}: already {}",
order.status()
);
self.clear_recon_tracking(&client_order_id, true);
return None;
}
if self.order_local_activity.within(
&client_order_id,
Duration::from_nanos(self.config.open_check_threshold_ns),
) {
return None;
}
let retries = self.recon_check_retries.entry(client_order_id).or_insert(0);
*retries = retries.saturating_add(1);
if *retries < self.config.open_check_missing_retries {
log::debug!(
"Order {} not found at venue, retry {}/{}",
client_order_id,
retries,
self.config.open_check_missing_retries
);
return None;
}
Some(order)
}
fn resolve_missing_order(&mut self, client_order_id: ClientOrderId) -> Vec<OrderEventAny> {
let mut events = Vec::new();
let Some(order) = self.get_order(client_order_id) else {
return events;
};
if order.status().is_closed() {
log::debug!(
"Skipping missing-order resolution for {client_order_id}: already {}",
order.status()
);
self.clear_recon_tracking(&client_order_id, true);
return events;
}
if self.order_local_activity.within(
&client_order_id,
Duration::from_nanos(self.config.open_check_threshold_ns),
) {
log::debug!(
"Deferring missing-order resolution for {client_order_id}: recent local activity"
);
return events;
}
let retries = self
.recon_check_retries
.get(&client_order_id)
.copied()
.unwrap_or_default();
let ts_now = self.clock.borrow().timestamp_ns();
match order.status() {
OrderStatus::Accepted | OrderStatus::Submitted => {
log::warn!(
"Order {client_order_id} not found at venue after {retries} retries and a targeted query, marking as REJECTED"
);
if let Some(rejected) =
create_reconciliation_rejected(&order, Some("NOT_FOUND_AT_VENUE"), ts_now)
{
events.push(rejected);
}
}
OrderStatus::PartiallyFilled => {
log::warn!(
"Order {client_order_id} not found at venue after {retries} retries and a targeted query, marking as CANCELED"
);
events.push(OrderEventAny::Canceled(OrderCanceled::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
client_order_id,
UUID4::new(),
ts_now,
ts_now,
true,
order.venue_order_id(),
order.account_id(),
)));
}
OrderStatus::PendingUpdate | OrderStatus::PendingCancel => {
log::debug!(
"Deferring resolution for {client_order_id}: still inflight as {}",
order.status()
);
self.recon_check_retries.shift_remove(&client_order_id);
if let Some(check) = self.inflight_checks.get_mut(&client_order_id) {
check.retry_count = 0;
check.last_query_at = Some(dst::time::Instant::now());
}
self.order_query_recency.mark(client_order_id);
return events;
}
status => {
log::warn!(
"Skipping missing-order resolution for {client_order_id}: unexpected status {status}"
);
}
}
self.clear_recon_tracking(&client_order_id, true);
events
}
fn check_position_discrepancy(
&mut self,
key: InstrumentAccountKey,
venue_reports: &[PositionStatusReport],
) -> Option<Vec<OrderEventAny>> {
let (instrument_id, account_id) = key;
let cached_positions = {
let cache = self.cache.borrow();
cache
.positions_open(None, Some(&instrument_id), None, Some(&account_id), None)
.into_iter()
.map(|position| (*position).clone())
.collect::<Vec<_>>()
};
let (cached_signed_qty, cached_long_qty, cached_short_qty) = Self::position_qty_aggregates(
cached_positions.iter().map(Position::signed_decimal_qty),
);
let (venue_signed_qty, venue_long_qty, venue_short_qty) = Self::position_qty_aggregates(
venue_reports.iter().map(|report| report.signed_decimal_qty),
);
let nonflat_count = venue_reports
.iter()
.filter(|report| report.signed_decimal_qty != Decimal::ZERO)
.count();
let venue_report = venue_reports
.iter()
.find(|report| report.signed_decimal_qty != Decimal::ZERO)
.or_else(|| venue_reports.last());
let tolerance = self.position_reconciliation_tolerance(account_id);
let venue_has_side_reports = venue_reports.iter().any(PositionStatusReport::is_long)
&& venue_reports.iter().any(PositionStatusReport::is_short);
let net_qty_matches = (cached_signed_qty - venue_signed_qty).abs() <= tolerance;
let side_qty_matches = (cached_long_qty - venue_long_qty).abs() <= tolerance
&& (cached_short_qty - venue_short_qty).abs() <= tolerance;
if net_qty_matches && (!venue_has_side_reports || side_qty_matches) {
self.position_reconciliation_states.shift_remove(&key);
return None;
}
let ts_now = self.clock.borrow().timestamp_ns();
if self.position_local_activity.within(
&key,
Duration::from_nanos(self.config.position_check_threshold_ns),
) {
log::debug!(
"Skipping position reconciliation for {instrument_id}: recent activity within threshold"
);
return None;
}
let report_shape = if nonflat_count > 1 || venue_has_side_reports {
PositionReportShape::MultiLeg
} else {
PositionReportShape::Unambiguous
};
let retries = self
.position_reconciliation_states
.get(&key)
.filter(|state| state.report_shape == report_shape)
.map_or(0, |state| state.retries);
if retries >= self.config.position_check_retries {
return None;
}
if report_shape == PositionReportShape::MultiLeg {
let new_retries = retries + 1;
self.set_position_reconciliation_retries(key, report_shape, new_retries);
log::warn!(
"Deferring position reconciliation for {instrument_id}/{account_id}: venue reports have ambiguous side aggregates (cached net={cached_signed_qty}, long={cached_long_qty}, short={cached_short_qty}; venue net={venue_signed_qty}, long={venue_long_qty}, short={venue_short_qty})"
);
if new_retries >= self.config.position_check_retries {
log::error!(
"Position discrepancy for {instrument_id}/{account_id} unresolved after {} attempts; no further reconciliation attempts will be made for the current report shape",
self.config.position_check_retries,
);
}
return None;
}
log::warn!(
"Position discrepancy detected for {instrument_id}: cached_signed_qty={cached_signed_qty}, venue_signed_qty={venue_signed_qty}"
);
let Some(instrument) = self.cache.borrow().instrument(&instrument_id).cloned() else {
log::debug!("Cannot reconcile position for {instrument_id}: instrument not in cache");
let new_retries = retries + 1;
self.set_position_reconciliation_retries(key, report_shape, new_retries);
if new_retries >= self.config.position_check_retries {
log::error!(
"Position discrepancy for {instrument_id} unresolved after {} attempts \
(cached_qty={cached_signed_qty}, venue_qty={venue_signed_qty}); \
no further reconciliation attempts will be made for the current report shape",
self.config.position_check_retries,
);
}
return None;
};
let cached_avg_px = Self::positions_avg_px(&cached_positions);
let venue_avg_px = venue_report.and_then(|r| r.avg_px_open);
let crosses_zero = (cached_signed_qty > Decimal::ZERO && venue_signed_qty < Decimal::ZERO)
|| (cached_signed_qty < Decimal::ZERO && venue_signed_qty > Decimal::ZERO);
let result = if crosses_zero {
let venue_ts_last = venue_report.map_or(ts_now, |r| r.ts_last);
self.reconcile_cross_zero_position(
&instrument,
account_id,
instrument_id,
cached_signed_qty,
cached_avg_px,
venue_signed_qty,
venue_avg_px,
ts_now,
venue_ts_last,
)
} else {
let qty_diff = venue_signed_qty - cached_signed_qty;
let order_side = if qty_diff > Decimal::ZERO {
OrderSide::Buy
} else {
OrderSide::Sell
};
let reconciliation_px = calculate_reconciliation_price(
cached_signed_qty,
cached_avg_px,
venue_signed_qty,
venue_avg_px,
);
match reconciliation_px.or(venue_avg_px).or(cached_avg_px) {
Some(fill_px) => {
let fill_qty = qty_diff.abs();
let venue_position_id = venue_report.and_then(|r| r.venue_position_id);
let venue_ts_last = venue_report.map_or(ts_now, |r| r.ts_last);
Quantity::from_decimal_dp(fill_qty, instrument.size_precision())
.ok()
.map(|order_qty| {
let fill_price =
Price::from_decimal_dp(fill_px, instrument.price_precision()).ok();
let venue_order_id = create_position_reconciliation_venue_order_id(
account_id,
instrument_id,
order_side,
OrderType::Market,
order_qty,
fill_price,
venue_position_id,
None,
venue_ts_last,
);
OrderStatusReport::new(
account_id,
instrument_id,
None,
venue_order_id,
order_side,
OrderType::Market,
TimeInForce::Gtc,
OrderStatus::Filled,
order_qty,
order_qty,
ts_now,
ts_now,
ts_now,
None,
)
.with_avg_px(fill_px)
})
.map(|order_report| {
log::info!(
color = LogColor::Blue as u8;
"Generating synthetic fill for position reconciliation {instrument_id}: side={order_side:?}, qty={}, px={fill_px}", qty_diff.abs(),
);
let (events, _) = self.handle_external_order(
&order_report,
account_id,
&instrument,
&[],
true,
None,
);
events
})
}
None => None,
}
};
if result.is_none() || result.as_ref().is_some_and(|e| e.is_empty()) {
let new_retries = retries + 1;
self.set_position_reconciliation_retries(key, report_shape, new_retries);
if new_retries >= self.config.position_check_retries {
log::error!(
"Position discrepancy for {} unresolved after {} attempts \
(cached_qty={}, venue_qty={}); \
no further reconciliation attempts will be made for the current report shape",
instrument_id,
self.config.position_check_retries,
cached_signed_qty,
venue_signed_qty,
);
}
} else {
self.position_reconciliation_states.shift_remove(&key);
}
result
}
fn set_position_reconciliation_retries(
&mut self,
key: InstrumentAccountKey,
report_shape: PositionReportShape,
retries: u32,
) {
self.position_reconciliation_states.insert(
key,
PositionReconciliationState {
report_shape,
retries,
},
);
}
fn position_qty_aggregates(
signed_quantities: impl Iterator<Item = Decimal>,
) -> (Decimal, Decimal, Decimal) {
signed_quantities.fold(
(Decimal::ZERO, Decimal::ZERO, Decimal::ZERO),
|(net, long, short), qty| {
if qty > Decimal::ZERO {
(net + qty, long + qty, short)
} else {
(net + qty, long, short + qty.abs())
}
},
)
}
#[expect(clippy::too_many_arguments)]
fn reconcile_cross_zero_position(
&self,
instrument: &InstrumentAny,
account_id: AccountId,
instrument_id: InstrumentId,
cached_signed_qty: Decimal,
cached_avg_px: Option<Decimal>,
venue_signed_qty: Decimal,
venue_avg_px: Option<Decimal>,
ts_now: UnixNanos,
venue_ts_last: UnixNanos,
) -> Option<Vec<OrderEventAny>> {
log::info!(
color = LogColor::Blue as u8;
"Position crosses zero for {instrument_id}: cached={cached_signed_qty}, venue={venue_signed_qty}. Splitting into two fills",
);
let close_qty = cached_signed_qty.abs();
let close_side = if cached_signed_qty < Decimal::ZERO {
OrderSide::Buy } else {
OrderSide::Sell };
let open_qty = venue_signed_qty.abs();
let open_side = if venue_signed_qty > Decimal::ZERO {
OrderSide::Buy } else {
OrderSide::Sell };
let Some(close_px) = cached_avg_px else {
log::warn!("Cannot close position for {instrument_id}: no cached average price");
return None;
};
let open_report = match venue_avg_px {
Some(open_px) => Some((
build_cross_zero_leg_report(
instrument,
account_id,
instrument_id,
open_side,
open_qty,
open_px,
"OPEN",
ts_now,
venue_ts_last,
)?,
open_px,
)),
None => None,
};
let close_report = build_cross_zero_leg_report(
instrument,
account_id,
instrument_id,
close_side,
close_qty,
close_px,
"CLOSE",
ts_now,
venue_ts_last,
)?;
log::info!(
color = LogColor::Blue as u8;
"Generating close fill for cross-zero {instrument_id}: side={close_side:?}, qty={close_qty}, px={close_px}",
);
let (close_events, _) =
self.handle_external_order(&close_report, account_id, instrument, &[], true, None);
let mut all_events = close_events;
if let Some((open_report, open_px)) = open_report {
log::info!(
color = LogColor::Blue as u8;
"Generating open fill for cross-zero {instrument_id}: side={open_side:?}, qty={open_qty}, px={open_px}",
);
let (open_events, _) =
self.handle_external_order(&open_report, account_id, instrument, &[], true, None);
all_events.extend(open_events);
} else {
log::warn!("Cannot open new position for {instrument_id}: no venue average price");
}
Some(all_events)
}
fn create_position_from_report(
&self,
report: &PositionStatusReport,
account_id: AccountId,
instrument: &InstrumentAny,
) -> Option<Vec<OrderEventAny>> {
let instrument_id = report.instrument_id;
let venue_signed_qty = report.signed_decimal_qty;
if venue_signed_qty == Decimal::ZERO {
return None;
}
let order_side = if venue_signed_qty > Decimal::ZERO {
OrderSide::Buy
} else {
OrderSide::Sell
};
let qty_abs = venue_signed_qty.abs();
let venue_avg_px = report.avg_px_open?;
let ts_now = self.clock.borrow().timestamp_ns();
let order_qty = Quantity::from_decimal_dp(qty_abs, instrument.size_precision()).ok()?;
let fill_price = Price::from_decimal_dp(venue_avg_px, instrument.price_precision()).ok();
let venue_order_id = create_position_reconciliation_venue_order_id(
account_id,
instrument_id,
order_side,
OrderType::Market,
order_qty,
fill_price,
report.venue_position_id,
None,
report.ts_last,
);
let mut order_report = OrderStatusReport::new(
account_id,
instrument_id,
None,
venue_order_id,
order_side,
OrderType::Market,
TimeInForce::Gtc,
OrderStatus::Filled,
order_qty,
order_qty,
ts_now,
ts_now,
ts_now,
None,
)
.with_avg_px(venue_avg_px);
if let Some(venue_position_id) = report.venue_position_id {
order_report = order_report.with_venue_position_id(venue_position_id);
}
log::info!(
color = LogColor::Blue as u8;
"Creating position from venue report for {instrument_id}: side={order_side:?}, qty={qty_abs}, avg_px={venue_avg_px}",
);
let (events, _) =
self.handle_external_order(&order_report, account_id, instrument, &[], true, None);
Some(events)
}
fn reconcile_position_report(
&self,
report: &PositionStatusReport,
account_id: AccountId,
instruments_with_unattributed_fills: &IndexSet<InstrumentId>,
positions_with_fills: &IndexSet<PositionId>,
) -> Option<Vec<OrderEventAny>> {
if report.venue_position_id.is_some() {
self.reconcile_position_report_hedging(
report,
account_id,
instruments_with_unattributed_fills,
positions_with_fills,
)
} else {
self.reconcile_position_report_netting(report, account_id)
}
}
fn reconcile_position_report_hedging(
&self,
report: &PositionStatusReport,
account_id: AccountId,
instruments_with_unattributed_fills: &IndexSet<InstrumentId>,
positions_with_fills: &IndexSet<PositionId>,
) -> Option<Vec<OrderEventAny>> {
let venue_position_id = report.venue_position_id?;
if positions_with_fills.contains(&venue_position_id) {
log::debug!(
"Skipping hedge position {venue_position_id} reconciliation: fills already in batch"
);
return None;
}
if instruments_with_unattributed_fills.contains(&report.instrument_id) {
log::debug!(
"Skipping hedge position {venue_position_id} reconciliation: unattributed fills in batch"
);
return None;
}
log::debug!(
"Reconciling HEDGE position for {}, venue_position_id={}",
report.instrument_id,
venue_position_id
);
let position = {
let cache = self.cache.borrow();
cache.position_owned(&venue_position_id)
};
match position {
Some(position) => {
let cached_signed_qty = position.signed_decimal_qty();
let venue_signed_qty = report.signed_decimal_qty;
if cached_signed_qty == venue_signed_qty {
log::debug!(
"Hedge position {venue_position_id} matches venue: qty={cached_signed_qty}"
);
return None;
}
if venue_signed_qty == Decimal::ZERO && cached_signed_qty == Decimal::ZERO {
return None;
}
if !self.config.generate_missing_orders {
log::error!(
"Cannot reconcile {} {}: position net qty {} != reported net qty {} \
and `generate_missing_orders` is disabled",
report.instrument_id,
venue_position_id,
cached_signed_qty,
venue_signed_qty
);
return None;
}
self.reconcile_hedge_position_discrepancy(
report,
account_id,
&position,
cached_signed_qty,
)
}
None => {
if report.signed_decimal_qty == Decimal::ZERO {
return None;
}
if !self.config.generate_missing_orders {
log::error!(
"Cannot reconcile position: {venue_position_id} not found and `generate_missing_orders` is disabled"
);
return None;
}
self.reconcile_missing_hedge_position(report, account_id)
}
}
}
fn reconcile_hedge_position_discrepancy(
&self,
report: &PositionStatusReport,
account_id: AccountId,
position: &Position,
cached_signed_qty: Decimal,
) -> Option<Vec<OrderEventAny>> {
let instrument = self.get_instrument(&report.instrument_id)?;
let venue_signed_qty = report.signed_decimal_qty;
let diff = (cached_signed_qty - venue_signed_qty).abs();
let diff_qty = Quantity::from_decimal_dp(diff, instrument.size_precision()).ok()?;
if diff_qty.is_zero() {
log::debug!(
"Difference quantity rounds to zero for {}, skipping",
instrument.id()
);
return None;
}
let venue_position_id = report.venue_position_id?;
log::warn!(
"Hedge position discrepancy for {} {}: cached={}, venue={}, generating reconciliation order",
report.instrument_id,
venue_position_id,
cached_signed_qty,
venue_signed_qty
);
let current_avg_px = if position.avg_px_open > 0.0 {
Decimal::from_str(&position.avg_px_open.to_string()).ok()
} else {
None
};
self.create_position_reconciliation_order(
report,
account_id,
&instrument,
cached_signed_qty,
diff_qty,
current_avg_px,
)
}
fn reconcile_missing_hedge_position(
&self,
report: &PositionStatusReport,
account_id: AccountId,
) -> Option<Vec<OrderEventAny>> {
let instrument = self.get_instrument(&report.instrument_id)?;
let venue_signed_qty = report.signed_decimal_qty;
let qty = venue_signed_qty.abs();
let diff_qty = Quantity::from_decimal_dp(qty, instrument.size_precision()).ok()?;
if diff_qty.is_zero() {
return None;
}
let venue_position_id = report.venue_position_id?;
log::warn!(
"Missing hedge position for {} {}: venue reports {}, generating reconciliation order",
report.instrument_id,
venue_position_id,
venue_signed_qty
);
self.create_position_reconciliation_order(
report,
account_id,
&instrument,
Decimal::ZERO,
diff_qty,
None,
)
}
fn reconcile_position_report_netting(
&self,
report: &PositionStatusReport,
account_id: AccountId,
) -> Option<Vec<OrderEventAny>> {
let instrument_id = report.instrument_id;
log::debug!("Reconciling NET position for {instrument_id}");
let instrument = self.get_instrument(&instrument_id)?;
let (cached_signed_qty, cached_avg_px) = {
let cache = self.cache.borrow();
let positions =
cache.positions_open(None, Some(&instrument_id), None, Some(&account_id), None);
if positions.is_empty() {
(Decimal::ZERO, None)
} else {
let mut total_signed_qty = Decimal::ZERO;
let mut total_value = Decimal::ZERO;
let mut total_qty = Decimal::ZERO;
for pos in positions {
total_signed_qty += pos.signed_decimal_qty();
let qty = pos.signed_decimal_qty().abs();
if pos.avg_px_open > 0.0
&& qty > Decimal::ZERO
&& let Ok(avg_px) = Decimal::from_str(&pos.avg_px_open.to_string())
{
total_value += avg_px * qty;
total_qty += qty;
}
}
let avg_px = if total_qty > Decimal::ZERO {
Some(total_value / total_qty)
} else {
None
};
(total_signed_qty, avg_px)
}
};
let venue_signed_qty = report.signed_decimal_qty;
log::debug!("venue_signed_qty={venue_signed_qty}, cached_signed_qty={cached_signed_qty}");
let tolerance = self.position_reconciliation_tolerance(account_id);
if (cached_signed_qty - venue_signed_qty).abs() <= tolerance {
log::debug!("Position quantities match for {instrument_id}, no reconciliation needed");
return None;
}
if !self.config.generate_missing_orders {
log::debug!(
"Discrepancy for {instrument_id} position when `generate_missing_orders` disabled, skipping"
);
return None;
}
let diff = (cached_signed_qty - venue_signed_qty).abs();
let diff_qty = Quantity::from_decimal_dp(diff, instrument.size_precision()).ok()?;
if diff_qty.is_zero() {
log::debug!(
"Difference quantity rounds to zero for {instrument_id}, skipping order generation"
);
return None;
}
let crosses_zero = cached_signed_qty != Decimal::ZERO
&& venue_signed_qty != Decimal::ZERO
&& ((cached_signed_qty > Decimal::ZERO && venue_signed_qty < Decimal::ZERO)
|| (cached_signed_qty < Decimal::ZERO && venue_signed_qty > Decimal::ZERO));
if crosses_zero {
let ts_now = self.clock.borrow().timestamp_ns();
return self.reconcile_cross_zero_position(
&instrument,
account_id,
instrument_id,
cached_signed_qty,
cached_avg_px,
venue_signed_qty,
report.avg_px_open,
ts_now,
report.ts_last,
);
}
if cached_signed_qty == Decimal::ZERO {
return self.create_position_from_report(report, account_id, &instrument);
}
self.create_position_reconciliation_order(
report,
account_id,
&instrument,
cached_signed_qty,
diff_qty,
cached_avg_px,
)
}
fn create_position_reconciliation_order(
&self,
report: &PositionStatusReport,
account_id: AccountId,
instrument: &InstrumentAny,
cached_signed_qty: Decimal,
diff_qty: Quantity,
current_avg_px: Option<Decimal>,
) -> Option<Vec<OrderEventAny>> {
let venue_signed_qty = report.signed_decimal_qty;
let instrument_id = report.instrument_id;
let order_side = if venue_signed_qty > cached_signed_qty {
OrderSide::Buy
} else {
OrderSide::Sell
};
let reconciliation_px = calculate_reconciliation_price(
cached_signed_qty,
current_avg_px,
venue_signed_qty,
report.avg_px_open,
);
let fill_px = reconciliation_px
.or(report.avg_px_open)
.or(current_avg_px)?;
let ts_now = self.clock.borrow().timestamp_ns();
let fill_price = Price::from_decimal_dp(fill_px, instrument.price_precision()).ok();
let venue_order_id = create_position_reconciliation_venue_order_id(
account_id,
instrument_id,
order_side,
OrderType::Market,
diff_qty,
fill_price,
report.venue_position_id,
None,
report.ts_last,
);
let mut order_report = OrderStatusReport::new(
account_id,
instrument_id,
None,
venue_order_id,
order_side,
OrderType::Market,
TimeInForce::Gtc,
OrderStatus::Filled,
diff_qty,
diff_qty,
ts_now,
ts_now,
ts_now,
None,
)
.with_avg_px(fill_px);
if let Some(venue_position_id) = report.venue_position_id {
order_report = order_report.with_venue_position_id(venue_position_id);
}
log::info!(
color = LogColor::Blue as u8;
"Generating reconciliation order for {instrument_id}: side={order_side:?}, qty={diff_qty}, px={fill_px}",
);
let (events, _) =
self.handle_external_order(&order_report, account_id, instrument, &[], true, None);
Some(events)
}
fn reconcile_order_report(
&self,
order: &OrderAny,
report: &OrderStatusReport,
instrument: Option<&InstrumentAny>,
) -> Option<OrderEventAny> {
let ts_now = self.clock.borrow().timestamp_ns();
reconcile_order_report(order, report, instrument, ts_now)
}
fn reconcile_order_with_fills(
&mut self,
order: &OrderAny,
report: &OrderStatusReport,
fills: &[&FillReport],
instrument: Option<&InstrumentAny>,
fill_queue: &mut ReconciliationFillQueue,
) -> Vec<OrderEventAny> {
let mut events = Vec::new();
let mut working = order.clone();
let mut sorted_fills: Vec<&FillReport> = fills.to_vec();
sorted_fills.sort_by_key(|f| f.ts_event);
let ts_now = self.clock.borrow().timestamp_ns();
if matches!(
report.order_status,
OrderStatus::Canceled | OrderStatus::Expired
) && report.ts_triggered.is_some()
&& working.status() != OrderStatus::Triggered
&& TRIGGERABLE_ORDER_TYPES.contains(&working.order_type())
{
let triggered = create_reconciliation_triggered(&working, report, ts_now);
if working.apply(triggered.clone()).is_ok() {
events.push(triggered);
}
}
let requires_snapshot_projection = !sorted_fills.is_empty()
|| report.order_status == OrderStatus::Voided
|| report.filled_qty < working.filled_qty();
if !requires_snapshot_projection {
if let Some(event) = self.reconcile_order_report(&working, report, instrument) {
events.push(event);
}
return events;
}
for event in generate_reconciliation_order_pre_fill_events(&working, report, ts_now) {
if let Err(e) = working.apply(event.clone()) {
log::warn!(
"Cannot project reconciliation event for {}: {e}",
order.client_order_id()
);
return events;
}
events.push(event);
}
if let Some(inst) = instrument {
for fill in sorted_fills {
let Some((event, fill_key)) =
self.create_order_fill(&working, fill, inst, &fill_queue.pending_fill_keys)
else {
continue;
};
if let Err(e) = working.apply(event.clone()) {
if let OrderEventAny::Filled(fill) = &event
&& self.is_fill_applied(fill, fill_key)
{
self.processed_fills.mark(fill_key);
} else {
log::warn!(
"Cannot project reconciliation fill for {}: {e}",
order.client_order_id()
);
}
return events;
}
fill_queue.push(&mut events, event, fill_key);
}
}
for event in
generate_reconciliation_order_snapshot_events(&working, report, instrument, ts_now)
{
if let Err(e) = working.apply(event.clone()) {
log::warn!(
"Cannot project reconciliation snapshot event for {}: {e}",
order.client_order_id()
);
break;
}
events.push(event);
}
events
}
fn handle_external_order(
&self,
report: &OrderStatusReport,
account_id: AccountId,
instrument: &InstrumentAny,
fills: &[&FillReport],
is_synthetic: bool,
mut fill_queue: Option<&mut ReconciliationFillQueue>,
) -> (Vec<OrderEventAny>, Option<ExternalOrderMetadata>) {
let (strategy_id, tags) =
if let Some(claimed_strategy) = self.external_order_claims.get(&report.instrument_id) {
let order_id = report
.client_order_id
.map_or_else(|| report.venue_order_id.to_string(), |id| id.to_string());
log::info!(
color = LogColor::Blue as u8;
"External order {} for {} claimed by strategy {}",
order_id,
report.instrument_id,
claimed_strategy,
);
(*claimed_strategy, None)
} else {
let tag = if is_synthetic {
*TAG_RECONCILIATION
} else {
*TAG_VENUE
};
(StrategyId::from("EXTERNAL"), Some(vec![tag]))
};
if self.config.filter_unclaimed_external && !is_synthetic {
return (Vec::new(), None);
}
let client_order_id = report
.client_order_id
.unwrap_or_else(|| ClientOrderId::from(report.venue_order_id.as_str()));
if !report.quantity.is_positive() {
log::error!(
"Skipping external order {} ({}) for {}: non-positive quantity in report {:?}",
client_order_id,
report.venue_order_id,
report.instrument_id,
report,
);
return (Vec::new(), None);
}
let ts_now = self.clock.borrow().timestamp_ns();
let initialized = match OrderInitialized::new_checked(
self.config.trader_id,
strategy_id,
report.instrument_id,
client_order_id,
report.order_side,
report.order_type,
report.quantity,
report.time_in_force,
report.post_only,
report.reduce_only,
false, true, UUID4::new(),
ts_now,
ts_now,
report.price,
report.activation_price,
report.trigger_price,
report.trigger_type,
report.limit_offset,
report.trailing_offset,
Some(report.trailing_offset_type),
report.expire_time,
report.display_qty,
None, None, Some(report.contingency_type),
report.order_list_id,
report.linked_order_ids.clone(),
report.parent_order_id,
None, None, None, tags,
) {
Ok(initialized) => initialized,
Err(e) => {
log::error!("Failed to create order from report: {e}");
return (Vec::new(), None);
}
};
let initialized = OrderEventAny::Initialized(initialized);
let order = match OrderAny::from_events(vec![initialized.clone()]) {
Ok(order) => order,
Err(e) => {
log::error!("Failed to create order from report: {e}");
return (Vec::new(), None);
}
};
{
let mut cache = self.cache.borrow_mut();
if let Err(e) = cache.add_order(order.clone(), None, None, false) {
match cache.order(&client_order_id) {
Some(existing) if is_synthetic && existing.is_closed() => {
log::debug!(
"Skipping synthetic reconciliation order {client_order_id} for {}: \
replay deduped (cached status={:?})",
report.instrument_id,
existing.status(),
);
}
Some(existing) if is_synthetic => {
log::warn!(
"Synthetic reconciliation order {client_order_id} for {} exists in \
cache in non-terminal state {:?}; fill not regenerated",
report.instrument_id,
existing.status(),
);
}
_ => {
log::error!("Failed to add external order to cache: {e}");
}
}
return (Vec::new(), None);
}
if let Err(e) =
cache.add_venue_order_id(&client_order_id, &report.venue_order_id, false)
{
log::warn!("Failed to add venue order ID index: {e}");
}
}
Self::publish_order_event(&initialized);
log::info!(
color = LogColor::Blue as u8;
"Created external order {} ({}) for {} [{}]",
client_order_id,
report.venue_order_id,
report.instrument_id,
report.order_status,
);
let ts_now = self.clock.borrow().timestamp_ns();
let mut order_events =
generate_external_order_status_events(&order, report, &account_id, instrument, ts_now);
if !fills.is_empty() {
let cached_order = self.get_order(client_order_id).unwrap();
let mut sorted_fills: Vec<&FillReport> = fills.to_vec();
sorted_fills.sort_by_key(|f| f.ts_event);
match report.order_status {
OrderStatus::Canceled
| OrderStatus::Expired
| OrderStatus::Filled
| OrderStatus::PartiallyFilled => {
let terminal_event = if order_events.last().is_some_and(|event| {
matches!(
event,
OrderEventAny::Canceled(_) | OrderEventAny::Expired(_),
)
}) {
order_events.pop()
} else {
None
};
if order_events
.last()
.is_some_and(|event| matches!(event, OrderEventAny::Filled(_)))
{
order_events.pop();
}
let mut real_fill_total = Decimal::ZERO;
for fill in &sorted_fills {
let fill_queue = fill_queue
.as_deref_mut()
.expect("real report fills require reconciliation queue state");
if let Some((fill_event, fill_key)) = self.create_order_fill(
&cached_order,
fill,
instrument,
&fill_queue.pending_fill_keys,
) {
real_fill_total += fill.last_qty.as_decimal();
fill_queue.push(&mut order_events, fill_event, fill_key);
}
}
let report_filled = report.filled_qty.as_decimal();
if real_fill_total < report_filled {
let diff_decimal = report_filled - real_fill_total;
if let Ok(diff) =
Quantity::from_decimal_dp(diff_decimal, instrument.size_precision())
&& let Some(inferred_fill) = create_inferred_fill_for_qty(
&cached_order,
report,
&account_id,
instrument,
diff,
ts_now,
None,
)
{
order_events.push(inferred_fill);
}
}
if let Some(event) = terminal_event {
order_events.push(event);
}
}
_ => {}
}
}
let metadata = ExternalOrderMetadata {
client_order_id,
venue_order_id: report.venue_order_id,
instrument_id: report.instrument_id,
strategy_id,
ts_init: ts_now,
};
(order_events, Some(metadata))
}
fn publish_order_event(event: &OrderEventAny) {
let topic = switchboard::get_event_order_topic(event.strategy_id());
msgbus::publish_order_event(topic, event);
}
fn adjust_mass_status_fills(
&self,
mass_status: &ExecutionMassStatus,
) -> (
IndexMap<VenueOrderId, OrderStatusReport>,
IndexMap<VenueOrderId, Vec<FillReport>>,
) {
let mut final_orders: IndexMap<VenueOrderId, OrderStatusReport> =
mass_status.order_reports();
let mut final_fills: IndexMap<VenueOrderId, Vec<FillReport>> = mass_status.fill_reports();
let mut instruments_to_adjust = Vec::new();
for (instrument_id, position_reports) in mass_status.position_reports() {
if !self.should_reconcile_instrument(&instrument_id) {
log::debug!(
"Skipping fill adjustment for {instrument_id}: not in reconciliation_instrument_ids"
);
continue;
}
let is_hedge_mode = position_reports
.iter()
.any(|r| r.venue_position_id.is_some());
if is_hedge_mode {
log::debug!(
"Skipping fill adjustment for {instrument_id}: hedge mode (has venue_position_id)"
);
continue;
}
let has_retained_position = {
let cache = self.cache.borrow();
!cache
.positions_open(
None,
Some(&instrument_id),
None,
Some(&mass_status.account_id),
None,
)
.is_empty()
};
if has_retained_position {
log::debug!(
"Skipping fill adjustment for {instrument_id}: retained open position in cache"
);
continue;
}
if let Some(instrument) = self.get_instrument(&instrument_id) {
instruments_to_adjust.push(instrument);
} else {
log::debug!(
"Skipping fill adjustment for {instrument_id}: instrument not found in cache"
);
}
}
if instruments_to_adjust.is_empty() {
return (final_orders, final_fills);
}
log_info!(
"Adjusting fills for {} instrument(s) with position reports",
instruments_to_adjust.len(),
color = LogColor::Blue
);
for instrument in &instruments_to_adjust {
let instrument_id = instrument.id();
match process_mass_status_for_reconciliation(mass_status, instrument, None) {
Ok(result) => {
final_orders.retain(|_, order| order.instrument_id != instrument_id);
final_fills.retain(|_, fills| {
fills
.first()
.is_none_or(|f| f.instrument_id != instrument_id)
});
for (venue_order_id, order) in result.orders {
final_orders.insert(venue_order_id, order);
}
for (venue_order_id, fills) in result.fills {
final_fills.insert(venue_order_id, fills);
}
}
Err(e) => {
log::warn!("Failed to adjust fills for {instrument_id}: {e}");
}
}
}
log_info!(
"After adjustment: {} order(s), {} fill group(s)",
final_orders.len(),
final_fills.len(),
color = LogColor::Blue
);
(final_orders, final_fills)
}
fn deduplicate_order_reports<'a>(
reports: impl Iterator<Item = &'a OrderStatusReport>,
) -> IndexMap<VenueOrderId, &'a OrderStatusReport> {
let mut best_reports: IndexMap<VenueOrderId, &'a OrderStatusReport> = IndexMap::new();
for report in reports {
let dominated = best_reports
.get(&report.venue_order_id)
.is_some_and(|existing| Self::is_more_advanced(existing, report));
if !dominated {
best_reports.insert(report.venue_order_id, report);
}
}
best_reports
}
fn is_more_advanced(a: &OrderStatusReport, b: &OrderStatusReport) -> bool {
if a.filled_qty > b.filled_qty {
return true;
}
if a.filled_qty < b.filled_qty {
return false;
}
Self::status_priority(a.order_status) > Self::status_priority(b.order_status)
}
const fn status_priority(status: OrderStatus) -> u8 {
match status {
OrderStatus::Initialized | OrderStatus::Submitted | OrderStatus::Emulated => 0,
OrderStatus::Released | OrderStatus::Denied => 1,
OrderStatus::Accepted | OrderStatus::PendingUpdate | OrderStatus::PendingCancel => 2,
OrderStatus::Triggered => 3,
OrderStatus::PartiallyFilled => 4,
OrderStatus::Canceled | OrderStatus::Expired | OrderStatus::Rejected => 5,
OrderStatus::Filled | OrderStatus::Voided => 6,
}
}
fn is_exact_order_match(order: &OrderAny, report: &OrderStatusReport) -> bool {
order.status() == report.order_status
&& order.filled_qty() == report.filled_qty
&& !should_reconciliation_update(order, report)
}
fn is_fill_applied(&self, fill: &OrderFilled, fill_key: FillKey) -> bool {
self.get_order(fill.client_order_id)
.or_else(|| self.get_order_by_venue_order_id(fill.venue_order_id))
.is_some_and(|order| {
order.account_id() == Some(fill_key.0)
&& order.instrument_id() == fill_key.1
&& order.trade_ids().contains(&&fill_key.2)
})
}
fn create_order_fill(
&self,
order: &OrderAny,
fill: &FillReport,
instrument: &InstrumentAny,
pending_fill_keys: &IndexSet<FillKey>,
) -> Option<(OrderEventAny, FillKey)> {
let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
if self.processed_fills.contains_key(&fill_key) || pending_fill_keys.contains(&fill_key) {
return None;
}
let event = OrderEventAny::Filled(OrderFilled::new(
order.trader_id(),
order.strategy_id(),
order.instrument_id(),
order.client_order_id(),
fill.venue_order_id,
fill.account_id,
fill.trade_id,
fill.order_side,
order.order_type(),
fill.last_qty,
fill.last_px,
instrument.quote_currency(),
fill.liquidity_side,
fill.report_id,
fill.ts_event,
self.clock.borrow().timestamp_ns(),
false,
fill.venue_position_id,
Some(fill.commission),
None,
));
Some((event, fill_key))
}
}
pub(crate) async fn request_targeted_order_reports(
clients: &[&dyn ExecutionClient],
queries: Vec<TargetedOrderQuery>,
query_delay: Duration,
) -> Vec<TargetedOrderReportResult> {
let mut results = Vec::with_capacity(queries.len());
let mut request_count = 0usize;
for query in queries {
let mut report = None;
let mut coverage_complete = true;
for client_id in &query.responsible_clients {
let client_id = *client_id;
let Some(client) = clients
.iter()
.find(|client| client.client_id() == client_id)
else {
coverage_complete = false;
log::warn!(
"Cannot run targeted order status query for {}: execution client {client_id} is unavailable",
query.client_order_id,
);
continue;
};
if request_count > 0 && !query_delay.is_zero() {
dst::time::sleep(query_delay).await;
}
request_count += 1;
match client.generate_order_status_report(&query.command).await {
Ok(Some(candidate)) if targeted_report_matches(&query, &candidate) => {
report = Some(candidate);
break;
}
Ok(Some(candidate)) => {
coverage_complete = false;
log::warn!(
"Ignoring mismatched targeted order status report from {client_id} for {}: client_order_id={:?}, venue_order_id={}, instrument_id={}",
query.client_order_id,
candidate.client_order_id,
candidate.venue_order_id,
candidate.instrument_id,
);
}
Ok(None) => {}
Err(e) => {
coverage_complete = false;
log::warn!(
"Failed targeted order status query from {client_id} for {}: {e}",
query.client_order_id,
);
}
}
}
results.push(TargetedOrderReportResult {
client_order_id: query.client_order_id,
report,
coverage_complete,
});
}
results
}
fn targeted_report_matches(query: &TargetedOrderQuery, report: &OrderStatusReport) -> bool {
let instrument_matches = query
.command
.instrument_id
.is_none_or(|instrument_id| report.instrument_id == instrument_id);
let order_matches = report.client_order_id == Some(query.client_order_id)
|| query
.command
.venue_order_id
.is_some_and(|venue_order_id| report.venue_order_id == venue_order_id);
instrument_matches && order_matches
}
#[cfg(test)]
mod tests {
use nautilus_common::clock::TestClock;
use nautilus_core::datetime::NANOSECONDS_IN_SECOND;
use nautilus_execution::reconciliation::generate_reconciliation_order_events;
use nautilus_model::{
enums::{LiquiditySide, OmsType, PositionSideSpecified},
events::order::spec::{OrderPendingUpdateSpec, OrderUpdatedSpec},
instruments::{
Instrument,
stubs::{crypto_perpetual_ethusdt, xbtusd_bitmex},
},
orders::{OrderTestBuilder, stubs::TestOrderEventStubs},
types::Money,
};
use rstest::rstest;
use super::*;
#[rstest]
fn test_clear_recon_tracking_removes_targeted_query() {
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let mut manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default());
let client_order_id = ClientOrderId::from("O-TARGETED-CLEAR");
manager.targeted_order_queries.insert(client_order_id);
manager.clear_recon_tracking(&client_order_id, true);
assert!(manager.targeted_order_queries.is_empty());
}
#[rstest]
fn test_register_inflight_skips_filtered_order() {
let client_order_id = ClientOrderId::from("O-FILTERED-REGISTER");
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let mut manager = ExecutionManager::new(
clock,
cache,
ExecutionManagerConfig {
filtered_client_order_ids: IndexSet::from([client_order_id]),
..Default::default()
},
);
manager.register_inflight(client_order_id);
assert!(!manager.inflight_checks.contains_key(&client_order_id));
assert!(!manager.recon_check_retries.contains_key(&client_order_id));
}
#[rstest]
#[cfg_attr(
not(all(feature = "simulation", madsim)),
tokio::test(start_paused = true)
)]
#[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
async fn test_inflight_check_retires_order_filtered_after_registration() {
let client_order_id = ClientOrderId::from("O-FILTERED-LATE");
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let mut manager = ExecutionManager::new(
clock,
cache,
ExecutionManagerConfig {
inflight_threshold_ms: 100,
..Default::default()
},
);
manager.register_inflight(client_order_id);
manager
.config
.filtered_client_order_ids
.insert(client_order_id);
dst::time::sleep(Duration::from_millis(101)).await;
let first = manager.check_inflight_orders();
assert!(first.events.is_empty());
assert!(first.queries.is_empty());
assert!(!manager.inflight_checks.contains_key(&client_order_id));
assert!(!manager.recon_check_retries.contains_key(&client_order_id));
dst::time::sleep(Duration::from_millis(101)).await;
let second = manager.check_inflight_orders();
assert!(second.events.is_empty());
assert!(second.queries.is_empty());
assert!(!manager.inflight_checks.contains_key(&client_order_id));
}
#[rstest]
#[case(false, OrderStatus::PendingUpdate, true, true, true)]
#[case(false, OrderStatus::Accepted, false, true, true)]
#[case(false, OrderStatus::Canceled, false, true, false)]
#[case(true, OrderStatus::PendingCancel, true, true, true)]
#[case(true, OrderStatus::Accepted, false, true, true)]
#[case(true, OrderStatus::Filled, false, true, false)]
fn test_observe_order_status_report_tracking_matrix(
#[case] with_fills: bool,
#[case] status: OrderStatus,
#[case] expect_inflight: bool,
#[case] expect_activity: bool,
#[case] expect_last_query: bool,
) {
let client_order_id = ClientOrderId::from("O-STATUS-MATRIX");
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let mut manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default());
manager.register_inflight(client_order_id);
manager.order_query_recency.mark(client_order_id);
manager
.missing_order_coverage_warnings
.insert(client_order_id);
manager.unresolved_order_coverage.insert(client_order_id);
manager.targeted_order_queries.insert(client_order_id);
let order_report = OrderStatusReport::new(
AccountId::from("TEST-001"),
crypto_perpetual_ethusdt().id(),
Some(client_order_id),
VenueOrderId::from("V-STATUS-MATRIX"),
OrderSide::Buy,
OrderType::Limit,
TimeInForce::Gtc,
status,
Quantity::from("10.0"),
Quantity::from("0.0"),
UnixNanos::from(1_000),
UnixNanos::from(1_000),
UnixNanos::from(1_000),
None,
);
let report = if with_fills {
ExecutionReport::OrderWithFills(Box::new(order_report), Vec::new())
} else {
ExecutionReport::Order(Box::new(order_report))
};
manager.observe_execution_report(&report);
assert_eq!(
manager.inflight_checks.contains_key(&client_order_id),
expect_inflight,
);
assert_eq!(
manager.recon_check_retries.contains_key(&client_order_id),
expect_inflight,
);
assert_eq!(
manager.order_local_activity.contains_key(&client_order_id),
expect_activity,
);
assert_eq!(
manager.order_query_recency.contains_key(&client_order_id),
expect_last_query,
);
assert_eq!(
manager
.missing_order_coverage_warnings
.contains(&client_order_id),
expect_inflight,
);
assert_eq!(
manager.unresolved_order_coverage.contains(&client_order_id),
expect_inflight,
);
assert_eq!(
manager.targeted_order_queries.contains(&client_order_id),
expect_inflight,
);
}
#[rstest]
fn test_superseded_cancel_report_preserves_missing_order_grace() {
let client_order_id = ClientOrderId::from("O-CANCEL-REPLACE");
let old_venue_order_id = VenueOrderId::from("V-CANCEL-REPLACE-OLD");
let new_venue_order_id = VenueOrderId::from("V-CANCEL-REPLACE-NEW");
let account_id = AccountId::from("TEST-001");
let client_id = ClientId::from("TEST");
let instrument_id = crypto_perpetual_ethusdt().id();
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
insert_accepted_limit_order(
&cache,
client_order_id,
old_venue_order_id,
instrument_id,
client_id,
);
let order = cache.borrow().order_owned(&client_order_id).unwrap();
let pending_update = OrderPendingUpdateSpec::builder()
.trader_id(order.trader_id())
.strategy_id(order.strategy_id())
.instrument_id(order.instrument_id())
.client_order_id(client_order_id)
.account_id(account_id)
.venue_order_id(old_venue_order_id)
.build();
cache
.borrow_mut()
.update_order(&OrderEventAny::PendingUpdate(pending_update))
.unwrap();
let order = cache.borrow().order_owned(&client_order_id).unwrap();
let updated = OrderUpdatedSpec::builder()
.trader_id(order.trader_id())
.strategy_id(order.strategy_id())
.instrument_id(order.instrument_id())
.client_order_id(client_order_id)
.quantity(order.quantity())
.venue_order_id(new_venue_order_id)
.account_id(account_id)
.build();
cache
.borrow_mut()
.update_order(&OrderEventAny::Updated(updated))
.unwrap();
let mut manager = ExecutionManager::new(
clock,
cache.clone(),
ExecutionManagerConfig {
open_check_missing_retries: 1,
..Default::default()
},
);
manager.record_local_activity(client_order_id);
assert!(
manager
.prepare_missing_order_query(client_order_id)
.is_none()
);
let report = OrderStatusReport::new(
account_id,
instrument_id,
Some(client_order_id),
old_venue_order_id,
OrderSide::Buy,
OrderType::Limit,
TimeInForce::Gtc,
OrderStatus::Canceled,
Quantity::from("10.0"),
Quantity::from("0.0"),
UnixNanos::from(1_000),
UnixNanos::from(2_000),
UnixNanos::from(3_000),
None,
);
manager.observe_execution_report(&ExecutionReport::Order(Box::new(report.clone())));
let order = cache.borrow().order_owned(&client_order_id).unwrap();
let events =
generate_reconciliation_order_events(&order, &report, None, UnixNanos::from(1_000));
assert!(events.is_empty());
assert_eq!(order.status(), OrderStatus::Accepted);
assert_eq!(order.venue_order_id(), Some(new_venue_order_id));
assert!(manager.order_local_activity.contains_key(&client_order_id));
assert!(
manager
.prepare_missing_order_query(client_order_id)
.is_none()
);
assert_eq!(manager.recon_check_retry_count(&client_order_id), 0);
}
#[rstest]
#[cfg_attr(
not(all(feature = "simulation", madsim)),
tokio::test(start_paused = true)
)]
#[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
async fn test_prune_order_local_activity_uses_open_check_threshold() {
let old_id = ClientOrderId::from("O-ACTIVITY-OLD");
let fresh_id = ClientOrderId::from("O-ACTIVITY-FRESH");
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let mut manager = ExecutionManager::new(
clock,
cache,
ExecutionManagerConfig {
open_check_threshold_ns: 100_000_000,
..Default::default()
},
);
manager.record_local_activity(old_id);
dst::time::sleep(Duration::from_millis(101)).await;
manager.record_local_activity(fresh_id);
manager.prune_order_local_activity();
assert!(!manager.order_local_activity.contains_key(&old_id));
assert!(manager.order_local_activity.contains_key(&fresh_id));
}
#[rstest]
fn test_prepare_open_order_report_check_builds_bulk_command_with_config() {
let lookback_mins = 5_u64;
let lookback_ns = lookback_mins * 60 * NANOSECONDS_IN_SECOND;
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let mut manager = ExecutionManager::new(
clock.clone(),
cache.clone(),
ExecutionManagerConfig {
open_check_lookback_mins: Some(lookback_mins),
open_check_open_only: false,
reconciliation_instrument_ids: IndexSet::from([crypto_perpetual_ethusdt().id()]),
..Default::default()
},
);
let included_id = ClientOrderId::from("O-REPORT-001");
let excluded_id = ClientOrderId::from("O-REPORT-002");
let included_instrument_id = crypto_perpetual_ethusdt().id();
let excluded_instrument_id = xbtusd_bitmex().id();
cache
.borrow_mut()
.add_instrument(InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt()))
.unwrap();
cache
.borrow_mut()
.add_instrument(InstrumentAny::CryptoPerpetual(xbtusd_bitmex()))
.unwrap();
insert_accepted_limit_order(
&cache,
included_id,
VenueOrderId::from("V-REPORT-001"),
included_instrument_id,
ClientId::from("BINANCE"),
);
insert_accepted_limit_order(
&cache,
excluded_id,
VenueOrderId::from("V-REPORT-002"),
excluded_instrument_id,
ClientId::from("BITMEX"),
);
clock
.borrow_mut()
.advance_time(UnixNanos::from(lookback_ns * 2), true);
let ts_now = clock.borrow().timestamp_ns();
let command_id = UUID4::new();
let check = manager.prepare_open_order_report_check(command_id, &[]);
assert_eq!(check.command.command_id, command_id);
assert_eq!(check.command.ts_init, ts_now);
assert!(!check.command.open_only);
assert_eq!(check.command.instrument_id, None);
assert_eq!(
check.command.start,
Some(ts_now.saturating_sub_ns(lookback_ns))
);
assert_eq!(check.command.end, None);
assert_eq!(check.command.log_receipt_level, LogLevel::Debug);
assert_eq!(check.start, check.command.start);
assert_eq!(check.filtered_orders.len(), 1);
assert_eq!(check.filtered_orders[0].client_order_id(), included_id);
}
#[rstest]
fn test_prepare_position_report_check_builds_bulk_command_with_coverage() {
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let manager = ExecutionManager::new(
clock.clone(),
cache.clone(),
ExecutionManagerConfig {
reconciliation_instrument_ids: IndexSet::from([crypto_perpetual_ethusdt().id()]),
..Default::default()
},
);
let included_instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let excluded_instrument = InstrumentAny::CryptoPerpetual(xbtusd_bitmex());
cache
.borrow_mut()
.add_instrument(included_instrument.clone())
.unwrap();
cache
.borrow_mut()
.add_instrument(excluded_instrument.clone())
.unwrap();
let included_position = insert_open_position(
&cache,
&included_instrument,
PositionId::from("P-REPORT-001"),
OrderSide::Buy,
"5.0",
"3000.00",
);
insert_open_position(
&cache,
&excluded_instrument,
PositionId::from("P-REPORT-002"),
OrderSide::Buy,
"2.0",
"40000.00",
);
let ts_now = clock.borrow().timestamp_ns();
let command_id = UUID4::new();
let check = manager.prepare_position_report_check(command_id, &[]);
let key = (
included_position.instrument_id,
included_position.account_id,
);
assert_eq!(check.command.command_id, command_id);
assert_eq!(check.command.ts_init, ts_now);
assert_eq!(check.command.instrument_id, None);
assert_eq!(check.command.start, None);
assert_eq!(check.command.end, None);
assert_eq!(check.command.log_receipt_level, LogLevel::Debug);
assert_eq!(check.client_coverage.len(), 1);
assert!(check.client_coverage.contains_key(&key));
assert_eq!(check.activity_revisions.get(&key), Some(&0));
}
#[rstest]
#[cfg_attr(
not(all(feature = "simulation", madsim)),
tokio::test(start_paused = true)
)]
#[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
async fn test_position_report_check_defers_activity_recorded_during_delayed_request() {
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let mut manager = ExecutionManager::new(
clock,
cache.clone(),
ExecutionManagerConfig {
position_check_threshold_ns: 5_000_000_000,
..Default::default()
},
);
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let instrument_id = instrument.id();
let position = insert_open_position(
&cache,
&instrument,
PositionId::from("P-ACTIVITY-DURING-REQUEST"),
OrderSide::Buy,
"5.0",
"3000.00",
);
cache
.borrow_mut()
.add_instrument(instrument.clone())
.unwrap();
let account_id = position.account_id;
let check = manager.prepare_position_report_check(UUID4::new(), &[]);
let report = PositionStatusReport::new(
account_id,
instrument_id,
PositionSideSpecified::Long,
Quantity::from("5.0"),
UnixNanos::from(1_000_000),
UnixNanos::from(1_000_000),
None,
None,
Some(Decimal::from(3000)),
);
let closed_position = close_long_position(
position,
&instrument,
TradeId::from("T-ACTIVITY-DURING-REQUEST"),
);
cache
.borrow_mut()
.update_position(&closed_position)
.unwrap();
manager.record_position_activity(instrument_id, account_id);
dst::time::sleep(Duration::from_secs(6)).await;
let events = manager.reconcile_position_reports(
&check,
vec![report],
&IndexSet::new(),
&IndexSet::new(),
);
assert!(
!events.iter().any(|event| {
matches!(
event,
OrderEventAny::Filled(fill)
if fill.order_side == OrderSide::Buy
&& fill.last_qty == Quantity::from("5.0")
)
}),
"activity recorded after the request started must defer A's stale report",
);
}
#[rstest]
fn test_position_report_check_does_not_defer_activity_recorded_before_request() {
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let mut manager = ExecutionManager::new(
clock,
cache.clone(),
ExecutionManagerConfig {
position_check_threshold_ns: 0,
..Default::default()
},
);
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let instrument_id = instrument.id();
let position = insert_open_position(
&cache,
&instrument,
PositionId::from("P-ACTIVITY-BEFORE-REQUEST"),
OrderSide::Buy,
"5.0",
"3000.00",
);
cache.borrow_mut().add_instrument(instrument).unwrap();
let account_id = position.account_id;
manager.record_position_activity(instrument_id, account_id);
let check = manager.prepare_position_report_check(UUID4::new(), &[]);
let report = PositionStatusReport::new(
account_id,
instrument_id,
PositionSideSpecified::Long,
Quantity::from("10.0"),
UnixNanos::from(1_000_000),
UnixNanos::from(1_000_000),
None,
None,
Some(Decimal::from(3000)),
);
let events = manager.reconcile_position_reports(
&check,
vec![report],
&IndexSet::new(),
&IndexSet::new(),
);
assert!(events.iter().any(|event| {
matches!(
event,
OrderEventAny::Filled(fill)
if fill.order_side == OrderSide::Buy
&& fill.last_qty == Quantity::from("5.0")
)
}));
}
#[rstest]
fn test_mass_status_projects_companion_fill_before_void_correction() {
let clock = Rc::new(RefCell::new(TestClock::new()));
let cache = Rc::new(RefCell::new(Cache::default()));
let mut manager =
ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default());
let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
let client_order_id = ClientOrderId::from("O-MASS-VOID-001");
let venue_order_id = VenueOrderId::from("V-MASS-VOID-001");
let account_id = AccountId::from("TEST-001");
cache
.borrow_mut()
.add_instrument(instrument.clone())
.unwrap();
insert_accepted_limit_order(
&cache,
client_order_id,
venue_order_id,
instrument.id(),
ClientId::from("BINANCE"),
);
let order = cache.borrow().order_owned(&client_order_id).unwrap();
let initial_fill = TestOrderEventStubs::filled(
&order,
&instrument,
Some(TradeId::from("T-MASS-VOID-INITIAL")),
None,
Some(Price::from("100.0")),
Some(Quantity::from("6.0")),
Some(LiquiditySide::Taker),
None,
None,
Some(account_id),
);
cache.borrow_mut().update_order(&initial_fill).unwrap();
let order = cache.borrow().order_owned(&client_order_id).unwrap();
let report = OrderStatusReport::new(
account_id,
instrument.id(),
Some(client_order_id),
venue_order_id,
OrderSide::Buy,
OrderType::Limit,
TimeInForce::Gtc,
OrderStatus::Canceled,
Quantity::from("10.0"),
Quantity::from("5.0"),
UnixNanos::from(1_000),
UnixNanos::from(1_000),
UnixNanos::from(1_000),
None,
);
let companion_fill = FillReport::new(
account_id,
instrument.id(),
venue_order_id,
TradeId::from("T-MASS-VOID-COMPANION"),
OrderSide::Buy,
Quantity::from("1.0"),
Price::from("100.0"),
Money::zero(instrument.quote_currency()),
LiquiditySide::Taker,
Some(client_order_id),
None,
UnixNanos::from(900),
UnixNanos::from(1_000),
None,
);
let mut fill_queue = ReconciliationFillQueue::default();
let events = manager.reconcile_order_with_fills(
&order,
&report,
&[&companion_fill],
Some(&instrument),
&mut fill_queue,
);
let mut projected = order;
for event in &events {
projected.apply(event.clone()).unwrap();
}
assert!(matches!(events[0], OrderEventAny::Filled(_)));
assert_eq!(
events
.iter()
.filter(|event| matches!(event, OrderEventAny::FillVoided(_)))
.count(),
2
);
assert_eq!(projected.status(), OrderStatus::Canceled);
assert_eq!(projected.filled_qty(), Quantity::from("5.0"));
assert_eq!(projected.voided_qty(), Quantity::from("2.0"));
}
fn insert_accepted_limit_order(
cache: &Rc<RefCell<Cache>>,
client_order_id: ClientOrderId,
venue_order_id: VenueOrderId,
instrument_id: InstrumentId,
client_id: ClientId,
) {
let account_id = AccountId::from("TEST-001");
let order = OrderTestBuilder::new(OrderType::Limit)
.client_order_id(client_order_id)
.instrument_id(instrument_id)
.quantity(Quantity::from("10.0"))
.price(Price::from("100.0"))
.build();
let submitted = TestOrderEventStubs::submitted(&order, account_id);
cache
.borrow_mut()
.add_order(order, None, Some(client_id), false)
.unwrap();
let order = cache.borrow_mut().update_order(&submitted).unwrap();
let accepted = TestOrderEventStubs::accepted(&order, account_id, venue_order_id);
cache.borrow_mut().update_order(&accepted).unwrap();
}
fn insert_open_position(
cache: &Rc<RefCell<Cache>>,
instrument: &InstrumentAny,
position_id: PositionId,
side: OrderSide,
quantity: &str,
price: &str,
) -> Position {
let order = OrderTestBuilder::new(OrderType::Market)
.instrument_id(instrument.id())
.side(side)
.quantity(Quantity::from(quantity))
.build();
let fill = TestOrderEventStubs::filled(
&order,
instrument,
Some(TradeId::new("T-REPORT-001")),
Some(position_id),
Some(Price::from(price)),
Some(Quantity::from(quantity)),
None,
None,
None,
Some(AccountId::from("TEST-001")),
);
let order_filled: OrderFilled = fill.into();
let position = Position::new(instrument, order_filled);
cache
.borrow_mut()
.add_position(&position, OmsType::Hedging)
.unwrap();
position
}
fn close_long_position(
mut position: Position,
instrument: &InstrumentAny,
trade_id: TradeId,
) -> Position {
let order = OrderTestBuilder::new(OrderType::Market)
.instrument_id(instrument.id())
.side(OrderSide::Sell)
.quantity(position.quantity)
.build();
let fill = TestOrderEventStubs::filled(
&order,
instrument,
Some(trade_id),
Some(position.id),
Some(Price::from("3000.00")),
Some(position.quantity),
None,
None,
None,
Some(position.account_id),
);
let order_filled: OrderFilled = fill.into();
position.apply(&order_filled);
position
}
}