1#[cfg(feature = "node")]
22use std::collections::HashSet;
23use std::{cell::RefCell, fmt::Debug, rc::Rc, str::FromStr, sync::LazyLock, time::Duration};
24
25use indexmap::{IndexMap, IndexSet};
26use nautilus_common::{
27 cache::Cache,
28 clients::{DEFAULT_POSITION_RECONCILIATION_TOLERANCE, ExecutionClient},
29 clock::Clock,
30 enums::{LogColor, LogLevel},
31 live::dst,
32 log_info,
33 messages::{
34 ExecutionReport,
35 execution::{
36 QueryOrder, TradingCommand,
37 report::{
38 GenerateOrderStatusReport, GenerateOrderStatusReports,
39 GeneratePositionStatusReports,
40 },
41 },
42 },
43 msgbus::{self, MessagingSwitchboard, switchboard},
44};
45use nautilus_core::{
46 UUID4, UnixNanos,
47 datetime::{mins_to_nanos, mins_to_secs},
48};
49use nautilus_execution::{
50 engine::ExecutionEngine,
51 reconciliation::{
52 calculate_reconciliation_price, create_inferred_fill_for_qty,
53 create_position_reconciliation_venue_order_id, create_reconciliation_rejected,
54 create_reconciliation_triggered, generate_external_order_status_events_with_commission,
55 generate_reconciliation_order_pre_fill_events,
56 generate_reconciliation_order_snapshot_events_with_commission,
57 incremental_inferred_fill_price_and_liquidity, inferred_fill_price_and_liquidity,
58 process_mass_status_for_reconciliation, reconcile_order_report_with_commission,
59 should_reconciliation_update,
60 },
61};
62use nautilus_model::{
63 enums::{LiquiditySide, OmsType, OrderSide, OrderStatus, OrderType, TimeInForce},
64 events::{OrderCanceled, OrderEventAny, OrderFilled, OrderInitialized},
65 identifiers::{
66 AccountId, ClientId, ClientOrderId, InstrumentId, PositionId, StrategyId, TradeId,
67 TraderId, VenueOrderId,
68 },
69 instruments::{Instrument, InstrumentAny},
70 orders::{Order, OrderAny, TRIGGERABLE_ORDER_TYPES},
71 position::Position,
72 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
73 types::{Money, Price, Quantity},
74};
75use rust_decimal::Decimal;
76use ustr::Ustr;
77
78use super::recency::RecencyMap;
79
80static TAG_VENUE: LazyLock<Ustr> = LazyLock::new(|| Ustr::from("VENUE"));
82
83static TAG_RECONCILIATION: LazyLock<Ustr> = LazyLock::new(|| Ustr::from("RECONCILIATION"));
85
86pub type InstrumentAccountKey = (InstrumentId, AccountId);
92type AccountInstrumentKey = (AccountId, InstrumentId);
93type AccountInstrumentStrategyKey = (AccountId, InstrumentId, StrategyId);
94type FillKey = (AccountId, InstrumentId, TradeId);
95
96#[expect(clippy::too_many_arguments)]
97fn build_cross_zero_leg_report(
98 instrument: &InstrumentAny,
99 account_id: AccountId,
100 instrument_id: InstrumentId,
101 order_side: OrderSide,
102 quantity: Decimal,
103 avg_px: Decimal,
104 tag: &str,
105 ts_now: UnixNanos,
106 venue_ts_last: UnixNanos,
107) -> Option<OrderStatusReport> {
108 let order_qty = Quantity::from_decimal_dp(quantity, instrument.size_precision()).ok()?;
109 let fill_price = Price::from_decimal_dp(avg_px, instrument.price_precision()).ok();
110 let venue_order_id = create_position_reconciliation_venue_order_id(
111 account_id,
112 instrument_id,
113 order_side,
114 OrderType::Market,
115 order_qty,
116 fill_price,
117 None,
118 Some(tag),
119 venue_ts_last,
120 );
121
122 let report = OrderStatusReport::new(
123 account_id,
124 instrument_id,
125 None,
126 venue_order_id,
127 order_side,
128 OrderType::Market,
129 TimeInForce::Gtc,
130 OrderStatus::Filled,
131 order_qty,
132 order_qty,
133 ts_now,
134 ts_now,
135 ts_now,
136 None,
137 )
138 .with_avg_px(avg_px);
139
140 Some(report)
141}
142
143#[derive(Debug, Clone, PartialEq, Eq)]
145pub(crate) enum ReportClientCoverage {
146 Resolved(IndexSet<ClientId>),
147 Unresolved,
148}
149
150#[derive(Debug, Clone)]
152pub struct ExternalOrderMetadata {
153 pub client_order_id: ClientOrderId,
154 pub venue_order_id: VenueOrderId,
155 pub instrument_id: InstrumentId,
156 pub strategy_id: StrategyId,
157 pub ts_init: UnixNanos,
158}
159
160#[derive(Debug, Default)]
162pub struct ReconciliationResult {
163 pub events: Vec<OrderEventAny>,
165 pub external_orders: Vec<ExternalOrderMetadata>,
167}
168
169#[derive(Debug, Default)]
171pub struct InflightCheckResult {
172 pub events: Vec<OrderEventAny>,
174 pub queries: Vec<TradingCommand>,
176}
177
178#[derive(Debug, Default)]
179pub(crate) struct OpenOrderReconciliationResult {
180 pub events: Vec<OrderEventAny>,
181 pub targeted_queries: Vec<TargetedOrderQuery>,
182}
183
184#[derive(Debug, Clone)]
185pub(crate) struct TargetedOrderQuery {
186 client_order_id: ClientOrderId,
187 responsible_clients: IndexSet<ClientId>,
188 command: GenerateOrderStatusReport,
189}
190
191impl TargetedOrderQuery {
192 #[cfg(feature = "node")]
193 pub(crate) const fn client_order_id(&self) -> ClientOrderId {
194 self.client_order_id
195 }
196}
197
198#[derive(Debug)]
199pub(crate) struct TargetedOrderReportResult {
200 client_order_id: ClientOrderId,
201 client_id: Option<ClientId>,
202 report: Option<OrderStatusReport>,
203 coverage_complete: bool,
204}
205
206#[derive(Debug)]
207pub(crate) struct SourcedOrderStatusReport {
208 pub client_id: ClientId,
209 pub report: OrderStatusReport,
210}
211
212#[derive(Debug, Clone)]
214pub(crate) struct OpenOrderReportCheck {
215 pub command: GenerateOrderStatusReports,
216 pub filtered_orders: Vec<OrderAny>,
217 pub client_coverage: IndexMap<ClientOrderId, ReportClientCoverage>,
218 pub start: Option<UnixNanos>,
219}
220
221#[derive(Debug, Clone)]
223pub(crate) struct PositionReportCheck {
224 pub command: GeneratePositionStatusReports,
225 pub client_coverage: IndexMap<InstrumentAccountKey, ReportClientCoverage>,
226 pub activity_revisions: IndexMap<InstrumentAccountKey, u64>,
227}
228
229struct RetainedFillState {
230 fill_keys: IndexSet<(AccountId, InstrumentId, TradeId)>,
231 missing_order_ids: IndexSet<(AccountId, InstrumentId, ClientOrderId)>,
232 missing_venue_order_ids: IndexSet<(AccountId, InstrumentId, VenueOrderId)>,
233 netting_lifecycle_starts: IndexMap<AccountInstrumentStrategyKey, UnixNanos>,
234}
235
236struct HistoricalFillGroup {
237 venue_order_id: VenueOrderId,
238 account_id: AccountId,
239 instrument_id: InstrumentId,
240 strategy_id: StrategyId,
241 order_side: OrderSide,
242 quantity: Decimal,
243 reduce_only: bool,
244 ts_event: UnixNanos,
245 ts_last: UnixNanos,
246}
247
248#[derive(Default)]
249struct ReconciliationFillQueue {
250 pending_fill_keys: IndexSet<FillKey>,
251 event_fill_keys: IndexMap<UUID4, FillKey>,
252}
253
254impl ReconciliationFillQueue {
255 fn push(&mut self, events: &mut Vec<OrderEventAny>, event: OrderEventAny, fill_key: FillKey) {
256 let OrderEventAny::Filled(fill) = &event else {
257 unreachable!("reported fills always create filled events");
258 };
259
260 self.pending_fill_keys.insert(fill_key);
261 self.event_fill_keys.insert(fill.event_id, fill_key);
262 events.push(event);
263 }
264}
265
266#[expect(
268 clippy::struct_excessive_bools,
269 reason = "config flags mirror the live execution engine configuration surface"
270)]
271#[derive(Debug, Clone)]
272pub struct ExecutionManagerConfig {
273 pub trader_id: TraderId,
275 pub reconciliation: bool,
277 pub lookback_mins: Option<u64>,
279 pub reconciliation_instrument_ids: IndexSet<InstrumentId>,
281 pub filter_unclaimed_external: bool,
283 pub filter_position_reports: bool,
285 pub filtered_client_order_ids: IndexSet<ClientOrderId>,
287 pub generate_missing_orders: bool,
289 pub inflight_check_interval_ms: u32,
291 pub inflight_threshold_ms: u64,
293 pub inflight_max_retries: u32,
295 pub open_check_interval_secs: Option<f64>,
297 pub open_check_lookback_mins: Option<u64>,
299 pub open_check_threshold_ns: u64,
301 pub open_check_missing_retries: u32,
303 pub open_check_open_only: bool,
305 pub max_single_order_queries_per_cycle: u32,
307 pub single_order_query_delay_ms: u32,
309 pub position_check_interval_secs: Option<f64>,
311 pub position_check_lookback_mins: u64,
313 pub position_check_threshold_ns: u64,
315 pub position_check_retries: u32,
317 pub purge_closed_orders_buffer_mins: Option<u32>,
319 pub purge_closed_positions_buffer_mins: Option<u32>,
321 pub purge_account_events_lookback_mins: Option<u32>,
323 pub purge_from_database: bool,
325}
326
327impl Default for ExecutionManagerConfig {
328 fn default() -> Self {
329 Self {
330 trader_id: TraderId::default(),
331 reconciliation: true,
332 lookback_mins: Some(60),
333 reconciliation_instrument_ids: IndexSet::new(),
334 filter_unclaimed_external: false,
335 filter_position_reports: false,
336 filtered_client_order_ids: IndexSet::new(),
337 generate_missing_orders: true,
338 inflight_check_interval_ms: 2_000,
339 inflight_threshold_ms: 5_000,
340 inflight_max_retries: 5,
341 open_check_interval_secs: None,
342 open_check_lookback_mins: Some(60),
343 open_check_threshold_ns: 5_000_000_000,
344 open_check_missing_retries: 5,
345 open_check_open_only: true,
346 max_single_order_queries_per_cycle: 5,
347 single_order_query_delay_ms: 100,
348 position_check_interval_secs: None,
349 position_check_lookback_mins: 60,
350 position_check_threshold_ns: 60_000_000_000,
351 position_check_retries: 3,
352 purge_closed_orders_buffer_mins: None,
353 purge_closed_positions_buffer_mins: None,
354 purge_account_events_lookback_mins: None,
355 purge_from_database: false,
356 }
357 }
358}
359
360impl ExecutionManagerConfig {
361 #[must_use]
363 pub fn with_trader_id(mut self, trader_id: TraderId) -> Self {
364 self.trader_id = trader_id;
365 self
366 }
367}
368
369#[derive(Debug, Clone)]
371struct InflightCheck {
372 #[allow(dead_code)]
373 pub client_order_id: ClientOrderId,
374 pub submitted_at: dst::time::Instant,
375 pub retry_count: u32,
376 pub last_query_at: Option<dst::time::Instant>,
379}
380
381#[derive(Clone, Copy, PartialEq, Eq)]
382enum PositionReportShape {
383 Unambiguous,
384 MultiLeg,
385}
386
387#[derive(Clone, Copy)]
388struct PositionReconciliationState {
389 report_shape: PositionReportShape,
390 retries: u32,
391}
392
393#[derive(Clone)]
416pub struct ExecutionManager {
417 clock: Rc<RefCell<dyn Clock>>,
418 cache: Rc<RefCell<Cache>>,
419 config: ExecutionManagerConfig,
420 inflight_checks: IndexMap<ClientOrderId, InflightCheck>,
421 external_order_claims: IndexMap<InstrumentId, StrategyId>,
422 processed_fills: RecencyMap<FillKey>,
423 recon_check_retries: IndexMap<ClientOrderId, u32>,
424 order_query_recency: RecencyMap<ClientOrderId>,
425 order_local_activity: RecencyMap<ClientOrderId>,
426 position_local_activity: RecencyMap<InstrumentAccountKey>,
428 position_local_activity_revisions: IndexMap<InstrumentAccountKey, u64>,
429 position_reconciliation_states: IndexMap<InstrumentAccountKey, PositionReconciliationState>,
430 position_reconciliation_tolerances: IndexMap<AccountId, Decimal>,
431 recent_fills_cache: RecencyMap<FillKey>,
432 missing_order_coverage_warnings: IndexSet<ClientOrderId>,
433 unresolved_order_coverage: IndexSet<ClientOrderId>,
434 targeted_order_queries: IndexSet<ClientOrderId>,
435}
436
437impl Debug for ExecutionManager {
438 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
439 f.debug_struct(stringify!(ExecutionManager))
440 .field("config", &self.config)
441 .field("inflight_checks", &self.inflight_checks)
442 .field("external_order_claims", &self.external_order_claims)
443 .field("processed_fills", &self.processed_fills)
444 .field("recon_check_retries", &self.recon_check_retries)
445 .finish_non_exhaustive()
446 }
447}
448
449impl ExecutionManager {
450 pub fn new(
452 clock: Rc<RefCell<dyn Clock>>,
453 cache: Rc<RefCell<Cache>>,
454 config: ExecutionManagerConfig,
455 ) -> Self {
456 Self {
457 clock,
458 cache,
459 config,
460 inflight_checks: IndexMap::new(),
461 external_order_claims: IndexMap::new(),
462 processed_fills: RecencyMap::default(),
463 recon_check_retries: IndexMap::new(),
464 order_query_recency: RecencyMap::default(),
465 order_local_activity: RecencyMap::default(),
466 position_local_activity: RecencyMap::default(),
467 position_local_activity_revisions: IndexMap::new(),
468 position_reconciliation_states: IndexMap::new(),
469 position_reconciliation_tolerances: IndexMap::new(),
470 recent_fills_cache: RecencyMap::default(),
471 missing_order_coverage_warnings: IndexSet::new(),
472 unresolved_order_coverage: IndexSet::new(),
473 targeted_order_queries: IndexSet::new(),
474 }
475 }
476
477 pub(crate) fn set_position_reconciliation_tolerance(
478 &mut self,
479 account_id: AccountId,
480 tolerance: Decimal,
481 ) {
482 let tolerance = if tolerance < Decimal::ZERO {
483 log::error!(
484 "Invalid negative position reconciliation tolerance {tolerance} for \
485 {account_id}; using the default"
486 );
487 DEFAULT_POSITION_RECONCILIATION_TOLERANCE
488 } else {
489 tolerance
490 };
491 self.position_reconciliation_tolerances
492 .insert(account_id, tolerance);
493 }
494
495 fn position_reconciliation_tolerance(&self, account_id: AccountId) -> Decimal {
496 self.position_reconciliation_tolerances
497 .get(&account_id)
498 .copied()
499 .unwrap_or(DEFAULT_POSITION_RECONCILIATION_TOLERANCE)
500 }
501
502 #[allow(unknown_lints, reason = "Clippy lint is unavailable on Rust 1.97")]
508 #[expect(
509 clippy::unused_async,
510 clippy::unused_async_trait_impl,
511 reason = "public reconciliation API stays async; live node and test callers await it"
512 )]
513 pub async fn reconcile_execution_mass_status(
514 &mut self,
515 mass_status: ExecutionMassStatus,
516 exec_engine: Rc<RefCell<ExecutionEngine>>,
517 ) -> ReconciliationResult {
518 if exec_engine
519 .borrow()
520 .get_client(&mass_status.client_id)
521 .is_none()
522 {
523 log::error!(
524 "Cannot reconcile ExecutionMassStatus from unknown client {}",
525 mass_status.client_id
526 );
527 return ReconciliationResult::default();
528 }
529
530 self.validate_mass_status_order_sources(&mass_status);
531
532 let raw_order_status_topic =
537 MessagingSwitchboard::reconciliation_raw_order_status_report_topic();
538
539 for report in mass_status.order_reports().values() {
540 msgbus::publish_any(raw_order_status_topic, report);
541 }
542
543 let raw_fill_topic = MessagingSwitchboard::reconciliation_raw_fill_report_topic();
544
545 for fills in mass_status.fill_reports().values() {
546 for fill in fills {
547 msgbus::publish_any(raw_fill_topic, fill);
548 }
549 }
550
551 let raw_position_topic =
552 MessagingSwitchboard::reconciliation_raw_position_status_report_topic();
553
554 for reports in mass_status.position_reports().values() {
555 for report in reports {
556 msgbus::publish_any(raw_position_topic, report);
557 }
558 }
559
560 if exec_engine
561 .borrow()
562 .get_client(&mass_status.client_id)
563 .is_none()
564 {
565 log::error!(
566 "Execution client {} disappeared while publishing raw mass status reports",
567 mass_status.client_id
568 );
569 return ReconciliationResult::default();
570 }
571
572 let venue = mass_status.venue;
573 let order_count = mass_status.order_reports().len();
574 let fill_count: usize = mass_status.fill_reports().values().map(Vec::len).sum();
575 let position_count = mass_status.position_reports().len();
576
577 log_info!(
578 "Reconciling ExecutionMassStatus for {venue}",
579 color = LogColor::Blue
580 );
581 log_info!(
582 "Received {order_count} order(s), {fill_count} fill(s), {position_count} position(s)",
583 color = LogColor::Blue
584 );
585
586 let retained_fill_state = self.retained_fill_state();
587 let reported_fill_keys: IndexSet<(AccountId, InstrumentId, TradeId)> = mass_status
588 .fill_reports()
589 .values()
590 .flatten()
591 .map(|fill| (fill.account_id, fill.instrument_id, fill.trade_id))
592 .collect();
593 let (adjusted_order_reports, adjusted_fill_reports) =
594 self.adjust_mass_status_fills(&mass_status);
595 let order_only_venue_order_ids = self.order_only_venue_order_ids(
596 &mass_status,
597 &adjusted_order_reports,
598 &adjusted_fill_reports,
599 &retained_fill_state,
600 );
601
602 let mut events = Vec::new();
603 let mut external_orders = Vec::new();
604 let mut orders_reconciled = 0usize;
605 let mut external_orders_created = 0usize;
606 let mut open_orders_initialized = 0usize;
607 let mut orders_skipped_no_instrument = 0usize;
608 let mut orders_skipped_duplicate = 0usize;
609 let mut fills_applied = 0usize;
610 let mut fill_queue = ReconciliationFillQueue::default();
611
612 let fill_reports = &adjusted_fill_reports;
613 let mut seen_fill_keys: IndexSet<FillKey> = IndexSet::new();
614
615 for fills in fill_reports.values() {
616 for fill in fills {
617 let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
618 if !seen_fill_keys.insert(fill_key) {
619 log::warn!(
620 "Duplicate trade_id {} for {} in mass status",
621 fill.trade_id,
622 fill.instrument_id
623 );
624 }
625 }
626 }
627
628 let order_reports = Self::deduplicate_order_reports(adjusted_order_reports.values());
630 let mut orders_skipped_filtered = 0usize;
631
632 for report in order_reports.values() {
633 if self.should_skip_order_report(report) {
634 orders_skipped_filtered += 1;
635 continue;
636 }
637
638 if let Some(client_order_id) = &report.client_order_id {
639 if let Some(cached_order) = self.get_order(*client_order_id)
640 && Self::is_exact_order_match(&cached_order, report)
641 {
642 log::debug!("Skipping order {client_order_id}: already in sync with venue");
643 orders_skipped_duplicate += 1;
644
645 if let Err(e) = self
647 .cache
648 .borrow_mut()
649 .index_venue_order_id(client_order_id, &report.venue_order_id)
650 {
651 log::warn!("Failed to index venue order ID: {e}");
652 }
653
654 continue;
655 }
656
657 if let Some(cached_order) = self.get_order(*client_order_id)
659 && cached_order.is_closed()
660 && cached_order
661 .tags()
662 .is_some_and(|tags| tags.contains(&*TAG_RECONCILIATION))
663 {
664 log::debug!(
665 "Skipping closed reconciliation order {client_order_id}: \
666 synthetic position adjustment from previous session",
667 );
668 orders_skipped_duplicate += 1;
669 continue;
670 }
671
672 if let Some(order) = self.get_order(*client_order_id) {
673 let instrument = self.get_instrument(&report.instrument_id);
674 log::info!(
675 color = LogColor::Blue as u8;
676 "Reconciling {} {} {} [{}] -> [{}]",
677 client_order_id,
678 report.venue_order_id,
679 report.instrument_id,
680 order.status(),
681 report.order_status,
682 );
683
684 let order_fills: Vec<&FillReport> = fill_reports
685 .get(&report.venue_order_id)
686 .map(|f| f.iter().collect())
687 .unwrap_or_default();
688 let engine_ref = exec_engine.borrow();
689 let commission_client = engine_ref.get_client(&mass_status.client_id);
690 let order_events = self.reconcile_order_with_fills(
691 &order,
692 report,
693 &order_fills,
694 instrument.as_ref(),
695 &mut fill_queue,
696 commission_client,
697 );
698 drop(engine_ref);
699
700 if !order_events.is_empty() {
701 orders_reconciled += 1;
702 fills_applied += order_events
703 .iter()
704 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
705 .count();
706 events.extend(order_events);
707 }
708
709 if let Err(e) = self
711 .cache
712 .borrow_mut()
713 .index_venue_order_id(client_order_id, &report.venue_order_id)
714 {
715 log::warn!("Failed to index venue order ID: {e}");
716 }
717 } else if let Some(order) = self.get_order_by_venue_order_id(report.venue_order_id)
718 {
719 let instrument = self.get_instrument(&report.instrument_id);
721
722 log::info!(
723 color = LogColor::Blue as u8;
724 "Reconciling {} (matched by venue_order_id {}) {} [{}] -> [{}]",
725 order.client_order_id(),
726 report.venue_order_id,
727 report.instrument_id,
728 order.status(),
729 report.order_status,
730 );
731
732 let order_fills: Vec<&FillReport> = fill_reports
733 .get(&report.venue_order_id)
734 .map(|f| f.iter().collect())
735 .unwrap_or_default();
736 let engine_ref = exec_engine.borrow();
737 let commission_client = engine_ref.get_client(&mass_status.client_id);
738 let order_events = self.reconcile_order_with_fills(
739 &order,
740 report,
741 &order_fills,
742 instrument.as_ref(),
743 &mut fill_queue,
744 commission_client,
745 );
746 drop(engine_ref);
747
748 if !order_events.is_empty() {
749 orders_reconciled += 1;
750 fills_applied += order_events
751 .iter()
752 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
753 .count();
754 events.extend(order_events);
755 }
756
757 if let Err(e) = self
758 .cache
759 .borrow_mut()
760 .index_venue_order_id(&order.client_order_id(), &report.venue_order_id)
761 {
762 log::warn!("Failed to index venue order ID: {e}");
763 }
764 } else if !self.config.filter_unclaimed_external {
765 if let Some(instrument) = self.get_instrument(&report.instrument_id) {
766 let order_fills: Vec<&FillReport> = fill_reports
767 .get(&report.venue_order_id)
768 .map(|f| f.iter().collect())
769 .unwrap_or_default();
770 let engine_ref = exec_engine.borrow();
771 let commission_client = engine_ref.get_client(&mass_status.client_id);
772 let (external_events, metadata) = self.handle_external_order(
773 report,
774 mass_status.account_id,
775 &instrument,
776 &order_fills,
777 false, Some(&mut fill_queue),
779 commission_client,
780 );
781 drop(engine_ref);
782
783 if !external_events.is_empty() {
784 external_orders_created += 1;
785 fills_applied += external_events
786 .iter()
787 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
788 .count();
789
790 if report.order_status.is_open() {
791 open_orders_initialized += 1;
792 }
793
794 events.extend(external_events);
795
796 if let Some(m) = metadata {
797 external_orders.push(m);
798 }
799 }
800 } else {
801 orders_skipped_no_instrument += 1;
802 }
803 }
804 } else if let Some(order) = self.get_order_by_venue_order_id(report.venue_order_id) {
805 let instrument = self.get_instrument(&report.instrument_id);
807 log::info!(
808 color = LogColor::Blue as u8;
809 "Reconciling {} (matched by venue_order_id {}) {} [{}] -> [{}]",
810 order.client_order_id(),
811 report.venue_order_id,
812 report.instrument_id,
813 order.status(),
814 report.order_status,
815 );
816
817 let order_fills: Vec<&FillReport> = fill_reports
818 .get(&report.venue_order_id)
819 .map(|f| f.iter().collect())
820 .unwrap_or_default();
821 let engine_ref = exec_engine.borrow();
822 let commission_client = engine_ref.get_client(&mass_status.client_id);
823 let order_events = self.reconcile_order_with_fills(
824 &order,
825 report,
826 &order_fills,
827 instrument.as_ref(),
828 &mut fill_queue,
829 commission_client,
830 );
831 drop(engine_ref);
832
833 if !order_events.is_empty() {
834 orders_reconciled += 1;
835 fills_applied += order_events
836 .iter()
837 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
838 .count();
839 events.extend(order_events);
840 }
841
842 if let Err(e) = self
843 .cache
844 .borrow_mut()
845 .index_venue_order_id(&order.client_order_id(), &report.venue_order_id)
846 {
847 log::warn!("Failed to index venue order ID: {e}");
848 }
849 } else if let Some(instrument) = self.get_instrument(&report.instrument_id) {
850 let is_synthetic = report.venue_order_id.as_str().starts_with("S-");
852
853 let order_fills: Vec<&FillReport> = fill_reports
854 .get(&report.venue_order_id)
855 .map(|f| f.iter().collect())
856 .unwrap_or_default();
857 let engine_ref = exec_engine.borrow();
858 let commission_client = engine_ref.get_client(&mass_status.client_id);
859 let (external_events, metadata) = self.handle_external_order(
860 report,
861 mass_status.account_id,
862 &instrument,
863 &order_fills,
864 is_synthetic,
865 Some(&mut fill_queue),
866 commission_client,
867 );
868 drop(engine_ref);
869
870 if !external_events.is_empty() {
871 external_orders_created += 1;
872 fills_applied += external_events
873 .iter()
874 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
875 .count();
876
877 if report.order_status.is_open() {
878 open_orders_initialized += 1;
879 }
880
881 events.extend(external_events);
882
883 if let Some(m) = metadata {
884 external_orders.push(m);
885 }
886 }
887 } else {
888 orders_skipped_no_instrument += 1;
889 }
890 }
891
892 let processed_venue_order_ids: IndexSet<VenueOrderId> =
894 order_reports.keys().copied().collect();
895
896 for (venue_order_id, fills) in fill_reports {
897 if processed_venue_order_ids.contains(venue_order_id) {
898 continue;
899 }
900
901 let Some(first_fill) = fills.first() else {
902 continue;
903 };
904
905 if !self.should_reconcile_instrument(&first_fill.instrument_id) {
906 log::debug!(
907 "Skipping orphan fills for {}: not in reconciliation_instrument_ids",
908 first_fill.instrument_id
909 );
910 continue;
911 }
912
913 if let Some(client_order_id) = &first_fill.client_order_id
915 && self
916 .config
917 .filtered_client_order_ids
918 .contains(client_order_id)
919 {
920 log::debug!(
921 "Skipping orphan fills for {client_order_id}: in filtered_client_order_ids"
922 );
923 continue;
924 }
925
926 let order = first_fill
927 .client_order_id
928 .as_ref()
929 .and_then(|id| self.get_order(*id))
930 .or_else(|| self.get_order_by_venue_order_id(*venue_order_id));
931
932 if let Some(ref order) = order
934 && self
935 .config
936 .filtered_client_order_ids
937 .contains(&order.client_order_id())
938 {
939 log::debug!(
940 "Skipping orphan fills for {}: in filtered_client_order_ids",
941 order.client_order_id()
942 );
943 continue;
944 }
945
946 if let Some(order) = order {
947 let instrument_id = order.instrument_id();
948 if let Some(instrument) = self.get_instrument(&instrument_id) {
949 let mut sorted_fills: Vec<&FillReport> = fills.iter().collect();
950 sorted_fills.sort_by_key(|f| f.ts_event);
951
952 for fill in sorted_fills {
953 if let Some((event, fill_key)) = self.create_order_fill(
954 &order,
955 fill,
956 &instrument,
957 &fill_queue.pending_fill_keys,
958 ) {
959 fills_applied += 1;
960 fill_queue.push(&mut events, event, fill_key);
961 }
962 }
963 }
964 }
965 }
966
967 events.sort_by_key(OrderEventAny::ts_event);
968
969 for event in &events {
970 if let OrderEventAny::Filled(fill) = event
971 && Self::should_project_reconciliation_fill(
972 fill,
973 &retained_fill_state,
974 &reported_fill_keys,
975 &order_only_venue_order_ids,
976 )
977 {
978 exec_engine.borrow_mut().project_reconciliation_fill(fill);
979 } else {
980 exec_engine.borrow_mut().process(event);
981 }
982
983 if let OrderEventAny::Filled(fill) = event
984 && let Some(fill_key) = fill_queue.event_fill_keys.get(&fill.event_id).copied()
985 && self.is_fill_applied(fill, fill_key)
986 {
987 self.processed_fills.mark(fill_key);
988 }
989 }
990
991 let mut positions_created = 0usize;
992
993 if !self.config.filter_position_reports {
994 let instruments_with_unattributed_fills: IndexSet<InstrumentId> = mass_status
997 .fill_reports()
998 .values()
999 .flatten()
1000 .filter(|f| f.venue_position_id.is_none())
1001 .map(|f| f.instrument_id)
1002 .chain(
1003 mass_status
1004 .order_reports()
1005 .values()
1006 .filter(|r| !r.filled_qty.is_zero() && r.venue_position_id.is_none())
1007 .map(|r| r.instrument_id),
1008 )
1009 .collect();
1010
1011 let positions_with_fills: IndexSet<PositionId> = mass_status
1012 .fill_reports()
1013 .values()
1014 .flatten()
1015 .filter_map(|f| f.venue_position_id)
1016 .chain(
1017 mass_status
1018 .order_reports()
1019 .values()
1020 .filter(|r| !r.filled_qty.is_zero())
1021 .filter_map(|r| r.venue_position_id),
1022 )
1023 .collect();
1024
1025 for (instrument_id, reports) in mass_status.position_reports() {
1026 if !self.should_reconcile_instrument(&instrument_id) {
1027 log::debug!(
1028 "Skipping position reports for {instrument_id}: not in reconciliation_instrument_ids"
1029 );
1030 continue;
1031 }
1032
1033 for report in reports {
1034 if let Some(position_events) = self.reconcile_position_report(
1035 &report,
1036 mass_status.account_id,
1037 &instruments_with_unattributed_fills,
1038 &positions_with_fills,
1039 ) {
1040 for event in position_events {
1041 exec_engine.borrow_mut().process(&event);
1042 events.push(event);
1043 }
1044 positions_created += 1;
1045 }
1046 }
1047 }
1048 }
1049
1050 if orders_skipped_no_instrument > 0 {
1051 log::warn!("{orders_skipped_no_instrument} orders skipped (instrument not in cache)");
1052 }
1053
1054 if orders_skipped_duplicate > 0 {
1055 log::debug!("{orders_skipped_duplicate} orders skipped (already in sync)");
1056 }
1057
1058 if orders_skipped_filtered > 0 {
1059 log::debug!("{orders_skipped_filtered} orders skipped (filtered by config)");
1060 }
1061
1062 log::info!(
1063 color = LogColor::Blue as u8;
1064 "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}",
1065 );
1066
1067 ReconciliationResult {
1068 events,
1069 external_orders,
1070 }
1071 }
1072
1073 fn should_project_reconciliation_fill(
1074 fill: &OrderFilled,
1075 retained_fill_state: &RetainedFillState,
1076 reported_fill_keys: &IndexSet<FillKey>,
1077 order_only_venue_order_ids: &IndexSet<VenueOrderId>,
1078 ) -> bool {
1079 let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
1080 if retained_fill_state.fill_keys.contains(&fill_key)
1081 || order_only_venue_order_ids.contains(&fill.venue_order_id)
1082 {
1083 return true;
1084 }
1085
1086 let order_missing = retained_fill_state.missing_order_ids.contains(&(
1087 fill.account_id,
1088 fill.instrument_id,
1089 fill.client_order_id,
1090 )) || retained_fill_state.missing_venue_order_ids.contains(&(
1091 fill.account_id,
1092 fill.instrument_id,
1093 fill.venue_order_id,
1094 ));
1095
1096 if order_missing && !reported_fill_keys.contains(&fill_key) {
1097 return true;
1098 }
1099
1100 retained_fill_state
1101 .netting_lifecycle_starts
1102 .get(&(fill.account_id, fill.instrument_id, fill.strategy_id))
1103 .is_some_and(|ts_opened| fill.ts_event < *ts_opened)
1104 }
1105
1106 fn retained_fill_state(&self) -> RetainedFillState {
1107 let cache = self.cache.borrow();
1108 let positions = cache.positions(None, None, None, None, None);
1109 let mut fill_keys = IndexSet::new();
1110 let mut missing_order_ids = IndexSet::new();
1111 let mut missing_venue_order_ids = IndexSet::new();
1112 let mut netting_lifecycle_starts = IndexMap::new();
1113
1114 for position in positions {
1115 for fill in &position.events {
1116 fill_keys.insert((position.account_id, position.instrument_id, fill.trade_id));
1117 if cache.order(&fill.client_order_id).is_none() {
1118 missing_order_ids.insert((
1119 position.account_id,
1120 position.instrument_id,
1121 fill.client_order_id,
1122 ));
1123 missing_venue_order_ids.insert((
1124 position.account_id,
1125 position.instrument_id,
1126 fill.venue_order_id,
1127 ));
1128 }
1129 }
1130
1131 if cache.oms_type(&position.id) == Some(OmsType::Netting) {
1132 netting_lifecycle_starts.insert(
1133 (
1134 position.account_id,
1135 position.instrument_id,
1136 position.strategy_id,
1137 ),
1138 position.ts_opened,
1139 );
1140 }
1141 }
1142
1143 RetainedFillState {
1144 fill_keys,
1145 missing_order_ids,
1146 missing_venue_order_ids,
1147 netting_lifecycle_starts,
1148 }
1149 }
1150
1151 fn order_only_venue_order_ids(
1152 &self,
1153 mass_status: &ExecutionMassStatus,
1154 order_reports: &IndexMap<VenueOrderId, OrderStatusReport>,
1155 fill_reports: &IndexMap<VenueOrderId, Vec<FillReport>>,
1156 retained_fill_state: &RetainedFillState,
1157 ) -> IndexSet<VenueOrderId> {
1158 if mass_status.lookback_start().is_none() {
1159 return IndexSet::new();
1160 }
1161
1162 let expected_quantities: IndexMap<AccountInstrumentKey, Decimal> =
1163 if mass_status.reports_complete() {
1164 mass_status
1165 .position_reports()
1166 .into_iter()
1167 .filter_map(|(instrument_id, reports)| {
1168 let [report] = reports.as_slice() else {
1169 return None;
1170 };
1171 report.venue_position_id.is_none().then_some((
1172 (report.account_id, instrument_id),
1173 report.signed_decimal_qty,
1174 ))
1175 })
1176 .collect()
1177 } else {
1178 IndexMap::new()
1179 };
1180 let candidate_instruments: IndexSet<InstrumentId> = order_reports
1181 .values()
1182 .filter(|report| !report.filled_qty.is_zero())
1183 .map(|report| report.instrument_id)
1184 .chain(
1185 fill_reports
1186 .values()
1187 .flatten()
1188 .map(|fill| fill.instrument_id),
1189 )
1190 .collect();
1191
1192 if candidate_instruments.is_empty() {
1193 return IndexSet::new();
1194 }
1195
1196 let mut venue_order_ids: IndexSet<VenueOrderId> = order_reports
1197 .iter()
1198 .filter(|(_, report)| {
1199 candidate_instruments.contains(&report.instrument_id)
1200 && !report.filled_qty.is_zero()
1201 })
1202 .map(|(venue_order_id, _)| *venue_order_id)
1203 .collect();
1204 venue_order_ids.extend(fill_reports.iter().filter_map(|(venue_order_id, fills)| {
1205 fills
1206 .first()
1207 .is_some_and(|fill| candidate_instruments.contains(&fill.instrument_id))
1208 .then_some(*venue_order_id)
1209 }));
1210
1211 let mut order_only = IndexSet::new();
1212 let mut groups = Vec::new();
1213
1214 for venue_order_id in venue_order_ids {
1215 let report = order_reports.get(&venue_order_id);
1216 let fills = fill_reports.get(&venue_order_id);
1217 if report.and_then(|report| report.venue_position_id).is_some()
1218 || fills.is_some_and(|fills| fills.iter().any(FillReport::has_venue_position_id))
1219 {
1220 continue;
1221 }
1222
1223 let cached_order = report
1224 .and_then(|report| report.client_order_id)
1225 .and_then(|client_order_id| self.get_order(client_order_id))
1226 .or_else(|| self.get_order_by_venue_order_id(venue_order_id));
1227 let account_id = report
1228 .map(|report| report.account_id)
1229 .or_else(|| fills.and_then(|fills| fills.first().map(|fill| fill.account_id)));
1230 let instrument_id = report
1231 .map(|report| report.instrument_id)
1232 .or_else(|| fills.and_then(|fills| fills.first().map(|fill| fill.instrument_id)));
1233 let order_side = report
1234 .map(|report| report.order_side)
1235 .or_else(|| fills.and_then(|fills| fills.first().map(|fill| fill.order_side)));
1236 let (Some(account_id), Some(instrument_id), Some(order_side)) =
1237 (account_id, instrument_id, order_side)
1238 else {
1239 order_only.insert(venue_order_id);
1240 continue;
1241 };
1242 let coherent_fills = fills.is_none_or(|fills| {
1243 fills.iter().all(|fill| {
1244 fill.account_id == account_id
1245 && fill.instrument_id == instrument_id
1246 && fill.order_side == order_side
1247 })
1248 });
1249 let coherent_cached_order = cached_order.as_ref().is_none_or(|order| {
1250 order.instrument_id() == instrument_id
1251 && order.order_side() == order_side
1252 && order.account_id().is_none_or(|id| id == account_id)
1253 });
1254
1255 if !coherent_fills
1256 || !coherent_cached_order
1257 || (report.is_none() && cached_order.is_none())
1258 {
1259 order_only.insert(venue_order_id);
1260 continue;
1261 }
1262 let strategy_id = cached_order.as_ref().map_or_else(
1263 || {
1264 self.external_order_claims
1265 .get(&instrument_id)
1266 .copied()
1267 .unwrap_or_else(|| StrategyId::from("EXTERNAL"))
1268 },
1269 Order::strategy_id,
1270 );
1271 let reduce_only = report.is_some_and(|report| report.reduce_only)
1272 || cached_order.as_ref().is_some_and(Order::is_reduce_only);
1273 let cached_filled_qty = cached_order
1274 .as_ref()
1275 .map_or(Decimal::ZERO, |order| order.filled_qty().as_decimal());
1276 let reported_fill_qty = fills.map_or(Decimal::ZERO, |fills| {
1277 fills.iter().map(|fill| fill.last_qty.as_decimal()).sum()
1278 });
1279 let unretained_fills: Vec<&FillReport> = fills
1280 .into_iter()
1281 .flatten()
1282 .filter(|fill| {
1283 !retained_fill_state.fill_keys.contains(&(
1284 fill.account_id,
1285 fill.instrument_id,
1286 fill.trade_id,
1287 ))
1288 })
1289 .collect();
1290 let unretained_fill_qty: Decimal = unretained_fills
1291 .iter()
1292 .map(|fill| fill.last_qty.as_decimal())
1293 .sum();
1294 let inferred_qty = report.map_or(Decimal::ZERO, |report| {
1295 (report.filled_qty.as_decimal() - cached_filled_qty - reported_fill_qty)
1296 .max(Decimal::ZERO)
1297 });
1298 let quantity = unretained_fill_qty + inferred_qty;
1299 if quantity.is_zero() {
1300 continue;
1301 }
1302 let inferred_ts = (!inferred_qty.is_zero())
1303 .then(|| report.map(|report| report.ts_last))
1304 .flatten();
1305 let ts_event = unretained_fills
1306 .iter()
1307 .map(|fill| fill.ts_event)
1308 .chain(inferred_ts)
1309 .min()
1310 .unwrap_or(mass_status.ts_init);
1311 let ts_last = unretained_fills
1312 .iter()
1313 .map(|fill| fill.ts_event)
1314 .chain(inferred_ts)
1315 .max()
1316 .unwrap_or(mass_status.ts_init);
1317 groups.push(HistoricalFillGroup {
1318 venue_order_id,
1319 account_id,
1320 instrument_id,
1321 strategy_id,
1322 order_side,
1323 quantity,
1324 reduce_only,
1325 ts_event,
1326 ts_last,
1327 });
1328 }
1329
1330 groups.sort_by_key(|group| group.ts_event);
1331 if !mass_status.reports_complete() {
1332 order_only.extend(groups.iter().map(|group| group.venue_order_id));
1333 log::error!(
1334 "Bounded reconciliation report set is incomplete; projecting {} historical order(s) without position or portfolio effects",
1335 order_only.len(),
1336 );
1337 return order_only;
1338 }
1339
1340 let mut quantities: IndexMap<AccountInstrumentStrategyKey, Option<Decimal>> =
1341 IndexMap::new();
1342 let mut group_ids: IndexMap<AccountInstrumentStrategyKey, Vec<VenueOrderId>> =
1343 IndexMap::new();
1344 let mut interval_ends: IndexMap<AccountInstrumentStrategyKey, UnixNanos> = IndexMap::new();
1345 let mut ambiguous_keys = IndexSet::new();
1346
1347 for group in &groups {
1348 let key = (group.account_id, group.instrument_id, group.strategy_id);
1349 if interval_ends
1350 .get(&key)
1351 .is_some_and(|end| group.ts_event <= *end)
1352 {
1353 ambiguous_keys.insert(key);
1354 }
1355 interval_ends
1356 .entry(key)
1357 .and_modify(|end| *end = (*end).max(group.ts_last))
1358 .or_insert(group.ts_last);
1359 }
1360
1361 if !ambiguous_keys.is_empty() {
1362 log::error!(
1363 "Bounded reconciliation contains interleaved order fills for {} position key(s); projecting their historical order state only",
1364 ambiguous_keys.len(),
1365 );
1366 }
1367
1368 for group in groups {
1369 let key = (group.account_id, group.instrument_id, group.strategy_id);
1370 group_ids.entry(key).or_default().push(group.venue_order_id);
1371 if ambiguous_keys.contains(&key) {
1372 order_only.insert(group.venue_order_id);
1373 continue;
1374 }
1375 let current_qty = quantities.entry(key).or_insert_with(|| {
1376 let cache = self.cache.borrow();
1377 let positions = cache.positions_open(
1378 None,
1379 Some(&group.instrument_id),
1380 Some(&group.strategy_id),
1381 Some(&group.account_id),
1382 None,
1383 );
1384
1385 if positions.len() > 1
1386 || positions.first().is_some_and(|position| {
1387 cache.oms_type(&position.id) != Some(OmsType::Netting)
1388 })
1389 {
1390 None
1391 } else {
1392 Some(
1393 positions
1394 .first()
1395 .map_or(Decimal::ZERO, |position| position.signed_decimal_qty()),
1396 )
1397 }
1398 });
1399 let Some(current_qty) = current_qty else {
1400 order_only.insert(group.venue_order_id);
1401 continue;
1402 };
1403 let signed_fill_qty = match group.order_side {
1404 OrderSide::Buy => group.quantity,
1405 OrderSide::Sell => -group.quantity,
1406 OrderSide::NoOrderSide => {
1407 order_only.insert(group.venue_order_id);
1408 continue;
1409 }
1410 };
1411 let reduces = !current_qty.is_zero()
1412 && current_qty.is_sign_negative() != signed_fill_qty.is_sign_negative()
1413 && group.quantity <= current_qty.abs();
1414 if group.reduce_only && !reduces {
1415 log::warn!(
1416 "Cannot apply bounded reduce-only order {} for {} without a coherent predecessor; projecting order state only",
1417 group.venue_order_id,
1418 group.instrument_id,
1419 );
1420 order_only.insert(group.venue_order_id);
1421 continue;
1422 }
1423 *current_qty += signed_fill_qty;
1424 }
1425
1426 let mut keys_by_position: IndexMap<
1427 AccountInstrumentKey,
1428 Vec<AccountInstrumentStrategyKey>,
1429 > = IndexMap::new();
1430
1431 for key in quantities.keys() {
1432 keys_by_position
1433 .entry((key.0, key.1))
1434 .or_default()
1435 .push(*key);
1436 }
1437
1438 for (position_key, keys) in keys_by_position {
1439 let expected_qty = expected_quantities.get(&position_key).copied();
1440 let matches_report = if expected_qty.is_some_and(|quantity| quantity.is_zero()) {
1441 keys.iter().all(|key| {
1442 quantities
1443 .get(key)
1444 .copied()
1445 .flatten()
1446 .is_some_and(|quantity| quantity.is_zero())
1447 })
1448 } else if let (Some(expected_qty), [key]) = (expected_qty, keys.as_slice()) {
1449 let cache = self.cache.borrow();
1450 let positions = cache.positions_open(
1451 None,
1452 Some(&position_key.1),
1453 None,
1454 Some(&position_key.0),
1455 None,
1456 );
1457 let cache_is_unambiguous = positions.len() <= 1
1458 && positions.first().is_none_or(|position| {
1459 position.strategy_id == key.2
1460 && cache.oms_type(&position.id) == Some(OmsType::Netting)
1461 });
1462 cache_is_unambiguous
1463 && quantities
1464 .get(key)
1465 .copied()
1466 .flatten()
1467 .is_some_and(|quantity| quantity == expected_qty)
1468 } else {
1469 false
1470 };
1471
1472 if matches_report {
1473 continue;
1474 }
1475
1476 let venue_order_ids: Vec<VenueOrderId> = keys
1477 .iter()
1478 .filter_map(|key| group_ids.get(key))
1479 .flatten()
1480 .copied()
1481 .collect();
1482 log::error!(
1483 "Bounded reconciliation does not explain the reported position for {}; projecting {} historical order(s) without position or portfolio effects",
1484 position_key.1,
1485 venue_order_ids.len(),
1486 );
1487 order_only.extend(venue_order_ids);
1488 }
1489
1490 order_only
1491 }
1492
1493 pub fn check_inflight_orders(&mut self) -> InflightCheckResult {
1499 let mut result = InflightCheckResult::default();
1500 let now = dst::time::Instant::now();
1501 let threshold = Duration::from_millis(self.config.inflight_threshold_ms);
1502
1503 let mut to_check = Vec::new();
1504
1505 for (client_order_id, check) in &self.inflight_checks {
1506 if now
1507 .checked_duration_since(check.submitted_at)
1508 .is_some_and(|elapsed| elapsed > threshold)
1509 {
1510 to_check.push(*client_order_id);
1511 }
1512 }
1513
1514 for client_order_id in to_check {
1515 if self
1516 .config
1517 .filtered_client_order_ids
1518 .contains(&client_order_id)
1519 {
1520 self.clear_recon_tracking(&client_order_id, true);
1521 continue;
1522 }
1523
1524 if self.targeted_order_queries.contains(&client_order_id) {
1525 continue;
1526 }
1527
1528 if let Some(check) = self.inflight_checks.get_mut(&client_order_id) {
1529 if let Some(last_query_at) = check.last_query_at
1530 && now
1531 .checked_duration_since(last_query_at)
1532 .is_none_or(|elapsed| elapsed < threshold)
1533 {
1534 continue;
1535 }
1536
1537 check.retry_count += 1;
1538 check.last_query_at = Some(now);
1539 self.order_query_recency.mark(client_order_id);
1540 self.recon_check_retries
1541 .insert(client_order_id, check.retry_count);
1542
1543 if check.retry_count >= self.config.inflight_max_retries {
1544 let ts_now = self.clock.borrow().timestamp_ns();
1545
1546 if let Some(order) = self.get_order(client_order_id) {
1547 match order.status() {
1548 OrderStatus::Submitted => {
1549 if let Some(event) = create_reconciliation_rejected(
1551 &order,
1552 Some("INFLIGHT_TIMEOUT"),
1553 ts_now,
1554 ) {
1555 result.events.push(event);
1556 }
1557 }
1558 OrderStatus::PendingUpdate | OrderStatus::PendingCancel => {
1559 let event = OrderEventAny::Canceled(OrderCanceled::new(
1561 order.trader_id(),
1562 order.strategy_id(),
1563 order.instrument_id(),
1564 order.client_order_id(),
1565 UUID4::new(),
1566 ts_now,
1567 ts_now,
1568 true, order.venue_order_id(),
1570 order.account_id(),
1571 ));
1572 result.events.push(event);
1573 }
1574 _ => {
1575 }
1577 }
1578 }
1579 self.clear_recon_tracking(&client_order_id, true);
1581 } else if let Some(order) = self.get_order(client_order_id) {
1582 let ts_now = self.clock.borrow().timestamp_ns();
1584 let client_id = self.cache.borrow().client_id(&client_order_id).copied();
1585 let query = TradingCommand::QueryOrder(QueryOrder::new(
1586 order.trader_id(),
1587 client_id,
1588 order.strategy_id(),
1589 order.instrument_id(),
1590 order.client_order_id(),
1591 order.venue_order_id(),
1592 UUID4::new(),
1593 ts_now,
1594 None,
1595 None, ));
1597 result.queries.push(query);
1598 }
1599 }
1600 }
1601
1602 result
1603 }
1604
1605 pub(crate) fn validate_mass_status_order_sources(&self, mass_status: &ExecutionMassStatus) {
1609 let cache = self.cache.borrow();
1610 let mut checked_client_order_ids = IndexSet::new();
1611 let mut missing_origins: Vec<ClientOrderId> = Vec::new();
1612 let mut mismatched_origins: Vec<(ClientOrderId, ClientId)> = Vec::new();
1613
1614 let mut validate_report_source =
1615 |direct_client_order_id: Option<ClientOrderId>, venue_order_id: VenueOrderId| {
1616 let direct_client_order_id = direct_client_order_id
1617 .filter(|client_order_id| cache.order_exists(client_order_id));
1618 let indexed_client_order_id = cache
1619 .client_order_id(&venue_order_id)
1620 .copied()
1621 .filter(|client_order_id| cache.order_exists(client_order_id));
1622
1623 for client_order_id in [direct_client_order_id, indexed_client_order_id]
1624 .into_iter()
1625 .flatten()
1626 .filter(|client_order_id| checked_client_order_ids.insert(*client_order_id))
1627 {
1628 match cache.client_id(&client_order_id) {
1629 Some(cached_client_id) if *cached_client_id == mass_status.client_id => {}
1630 Some(cached_client_id) => {
1631 mismatched_origins.push((client_order_id, *cached_client_id));
1632 }
1633 None => missing_origins.push(client_order_id),
1634 }
1635 }
1636 };
1637
1638 for report in mass_status.order_reports().values() {
1639 validate_report_source(report.client_order_id, report.venue_order_id);
1640 }
1641
1642 for fills in mass_status.fill_reports().values() {
1643 for fill in fills {
1644 validate_report_source(fill.client_order_id, fill.venue_order_id);
1645 }
1646 }
1647
1648 if !missing_origins.is_empty() {
1649 let samples = missing_origins
1650 .iter()
1651 .take(5)
1652 .map(ToString::to_string)
1653 .collect::<Vec<_>>()
1654 .join(", ");
1655
1656 log::warn!(
1657 "Found {} cached order(s) without an execution client origin ({}): \
1658 continuing reconciliation against mass status client {} for compatibility \
1659 with existing cache data",
1660 missing_origins.len(),
1661 samples,
1662 mass_status.client_id,
1663 );
1664 }
1665
1666 if !mismatched_origins.is_empty() {
1667 let samples = mismatched_origins
1668 .iter()
1669 .take(5)
1670 .map(|(client_order_id, cached)| format!("{client_order_id} -> {cached}"))
1671 .collect::<Vec<_>>()
1672 .join(", ");
1673
1674 log::warn!(
1675 "Found {} cached order(s) with an execution client origin conflicting with \
1676 mass status client {} ({}): continuing reconciliation for compatibility; \
1677 this conflict will become a startup error in a future release, verify cached \
1678 order ownership and execution client configuration",
1679 mismatched_origins.len(),
1680 mass_status.client_id,
1681 samples,
1682 );
1683 }
1684 }
1685
1686 fn filtered_open_orders_for_reconciliation(&self) -> Vec<OrderAny> {
1687 {
1688 let cache = self.cache.borrow();
1689 let mut orders = cache.orders_open(None, None, None, None, None);
1690 orders.extend(cache.orders_inflight(None, None, None, None, None));
1691 let mut seen_client_order_ids = IndexSet::new();
1692 orders.retain(|order| seen_client_order_ids.insert(order.client_order_id()));
1693
1694 if self.config.reconciliation_instrument_ids.is_empty() {
1695 orders.iter().map(|o| (*o).clone()).collect()
1696 } else {
1697 orders
1698 .iter()
1699 .filter(|o| {
1700 self.config
1701 .reconciliation_instrument_ids
1702 .contains(&o.instrument_id())
1703 })
1704 .map(|o| (*o).clone())
1705 .collect()
1706 }
1707 }
1708 }
1709
1710 fn open_position_keys_for_reconciliation(&self) -> IndexSet<InstrumentAccountKey> {
1711 let cache = self.cache.borrow();
1712 let positions = cache.positions_open(None, None, None, None, None);
1713 let mut position_keys = IndexSet::new();
1714
1715 for position in positions {
1716 if !self.should_reconcile_instrument(&position.instrument_id) {
1717 continue;
1718 }
1719
1720 position_keys.insert((position.instrument_id, position.account_id));
1721 }
1722
1723 position_keys
1724 }
1725
1726 pub(crate) fn prepare_open_order_report_check(
1728 &mut self,
1729 command_id: UUID4,
1730 clients: &[&dyn ExecutionClient],
1731 ) -> OpenOrderReportCheck {
1732 let filtered_orders = self.filtered_open_orders_for_reconciliation();
1733 let active_order_ids: IndexSet<ClientOrderId> =
1734 filtered_orders.iter().map(Order::client_order_id).collect();
1735 self.missing_order_coverage_warnings
1736 .retain(|client_order_id| active_order_ids.contains(client_order_id));
1737 self.unresolved_order_coverage
1738 .retain(|client_order_id| active_order_ids.contains(client_order_id));
1739
1740 let mut client_coverage = IndexMap::new();
1741
1742 for order in &filtered_orders {
1743 let client_order_id = order.client_order_id();
1744 let coverage = self.resolve_order_report_client_coverage(order, clients);
1745
1746 match &coverage {
1747 ReportClientCoverage::Resolved(_) => {
1748 if self
1749 .unresolved_order_coverage
1750 .shift_remove(&client_order_id)
1751 {
1752 self.missing_order_coverage_warnings
1753 .shift_remove(&client_order_id);
1754 }
1755 }
1756 ReportClientCoverage::Unresolved => {
1757 self.unresolved_order_coverage.insert(client_order_id);
1758 }
1759 }
1760
1761 client_coverage.insert(client_order_id, coverage);
1762 }
1763
1764 log::debug!(
1765 "Found {} order{} open in cache",
1766 filtered_orders.len(),
1767 if filtered_orders.len() == 1 { "" } else { "s" }
1768 );
1769
1770 let ts_now = self.clock.borrow().timestamp_ns();
1771 let start = self.config.open_check_lookback_mins.map(|mins| {
1772 let lookback_ns = mins_to_nanos(mins);
1773 ts_now.saturating_sub_ns(lookback_ns)
1774 });
1775
1776 let mut command = GenerateOrderStatusReports::new(
1777 command_id,
1778 ts_now,
1779 self.config.open_check_open_only,
1780 None,
1781 start,
1782 None,
1783 None,
1784 None,
1785 );
1786 command.log_receipt_level = LogLevel::Debug;
1787
1788 OpenOrderReportCheck {
1789 command,
1790 filtered_orders,
1791 client_coverage,
1792 start,
1793 }
1794 }
1795
1796 fn resolve_order_report_client_coverage(
1797 &self,
1798 order: &OrderAny,
1799 clients: &[&dyn ExecutionClient],
1800 ) -> ReportClientCoverage {
1801 if let Some(client_id) = self.cache.borrow().client_id(&order.client_order_id()) {
1802 return ReportClientCoverage::Resolved(IndexSet::from([*client_id]));
1803 }
1804
1805 if let Some(account_id) = order.account_id() {
1806 let account_clients = clients
1807 .iter()
1808 .filter(|client| client.account_id() == account_id)
1809 .map(|client| client.client_id())
1810 .collect::<IndexSet<_>>();
1811
1812 if !account_clients.is_empty() {
1813 return ReportClientCoverage::Resolved(account_clients);
1814 }
1815 }
1816
1817 let venue_clients = clients
1818 .iter()
1819 .filter(|client| client.handles_order_venue(order.instrument_id().venue))
1820 .map(|client| client.client_id())
1821 .collect::<IndexSet<_>>();
1822
1823 if venue_clients.is_empty() {
1824 ReportClientCoverage::Unresolved
1825 } else {
1826 ReportClientCoverage::Resolved(venue_clients)
1827 }
1828 }
1829
1830 pub fn check_open_order_queries(&mut self) -> Vec<TradingCommand> {
1832 self.check_open_order_queries_for_clients(None)
1833 }
1834
1835 pub(crate) fn check_open_order_queries_for_clients(
1836 &mut self,
1837 client_ids: Option<&IndexSet<ClientId>>,
1838 ) -> Vec<TradingCommand> {
1839 let now = dst::time::Instant::now();
1840 let query_delay = Duration::from_millis(u64::from(self.config.single_order_query_delay_ms));
1841 let query_limit = self.config.max_single_order_queries_per_cycle as usize;
1842
1843 if query_limit == 0 {
1844 return Vec::new();
1845 }
1846
1847 let mut filtered_orders = self.filtered_open_orders_for_reconciliation();
1848 filtered_orders.sort_by_key(|order| {
1849 let client_order_id = order.client_order_id();
1850 (
1851 self.order_query_recency.last_marked(&client_order_id),
1852 client_order_id,
1853 )
1854 });
1855
1856 let mut queries = Vec::new();
1857
1858 for order in filtered_orders {
1859 if queries.len() >= query_limit {
1860 break;
1861 }
1862
1863 let client_order_id = order.client_order_id();
1864 let client_id = self.cache.borrow().client_id(&client_order_id).copied();
1865
1866 if let Some(client_ids) = client_ids
1867 && !client_id.is_some_and(|client_id| client_ids.contains(&client_id))
1868 {
1869 continue;
1870 }
1871
1872 if self
1873 .config
1874 .filtered_client_order_ids
1875 .contains(&client_order_id)
1876 {
1877 continue;
1878 }
1879
1880 let threshold = Duration::from_nanos(self.config.open_check_threshold_ns);
1881 if let Some(elapsed) = self.order_local_activity.elapsed_at(&client_order_id, now)
1882 && elapsed < threshold
1883 {
1884 let elapsed_ms = elapsed.as_millis();
1885 let threshold_ms = threshold.as_millis();
1886 log::debug!(
1887 "Deferring open order query for {client_order_id}: recent local activity \
1888 ({elapsed_ms}ms < threshold={threshold_ms}ms)",
1889 );
1890 continue;
1891 }
1892
1893 if self
1894 .order_query_recency
1895 .within_at(&client_order_id, now, query_delay)
1896 {
1897 continue;
1898 }
1899
1900 self.order_query_recency.mark(client_order_id);
1901 let ts_now = self.clock.borrow().timestamp_ns();
1902
1903 let cmd = TradingCommand::QueryOrder(QueryOrder::new(
1904 order.trader_id(),
1905 client_id,
1906 order.strategy_id(),
1907 order.instrument_id(),
1908 client_order_id,
1909 order.venue_order_id(),
1910 UUID4::new(),
1911 ts_now,
1912 None,
1913 None,
1914 ));
1915 queries.push(cmd);
1916 }
1917
1918 queries
1919 }
1920
1921 pub async fn check_open_orders(
1931 &mut self,
1932 clients: &[&dyn ExecutionClient],
1933 ) -> Vec<OrderEventAny> {
1934 log::debug!("Checking order consistency between cached-state and venues");
1935
1936 let check = self.prepare_open_order_report_check(UUID4::new(), clients);
1937 let mut all_reports = Vec::new();
1938 let mut queried_clients = IndexSet::new();
1939 let mut failed_clients = IndexSet::new();
1940
1941 for client in clients {
1942 let client_id = client.client_id();
1943 queried_clients.insert(client_id);
1944
1945 match client.generate_order_status_reports(&check.command).await {
1946 Ok(reports) => {
1947 all_reports.extend(
1948 reports
1949 .into_iter()
1950 .map(|report| SourcedOrderStatusReport { client_id, report }),
1951 );
1952 }
1953 Err(e) => {
1954 failed_clients.insert(client_id);
1955 log::warn!(
1956 "Failed to query order reports from {}: {e}",
1957 client.client_id()
1958 );
1959 }
1960 }
1961 }
1962
1963 let result = self.reconcile_open_order_reports(
1964 &check,
1965 all_reports,
1966 &queried_clients,
1967 &failed_clients,
1968 clients,
1969 );
1970 let mut events = result.events;
1971
1972 if !result.targeted_queries.is_empty() {
1973 let query_delay =
1974 Duration::from_millis(u64::from(self.config.single_order_query_delay_ms));
1975 let query_results =
1976 request_targeted_order_reports(clients, result.targeted_queries, query_delay).await;
1977 events.extend(self.reconcile_targeted_order_reports(query_results, clients));
1978 }
1979
1980 events
1981 }
1982
1983 pub(crate) fn reconcile_open_order_reports(
1985 &mut self,
1986 check: &OpenOrderReportCheck,
1987 all_reports: Vec<SourcedOrderStatusReport>,
1988 queried_clients: &IndexSet<ClientId>,
1989 failed_clients: &IndexSet<ClientId>,
1990 clients: &[&dyn ExecutionClient],
1991 ) -> OpenOrderReconciliationResult {
1992 let mut venue_reported_ids = IndexSet::new();
1993
1994 for sourced in &all_reports {
1995 let report = &sourced.report;
1996 if let Some(client_order_id) = &report.client_order_id {
1997 venue_reported_ids.insert(*client_order_id);
1998 self.missing_order_coverage_warnings
1999 .shift_remove(client_order_id);
2000 self.recon_check_retries.shift_remove(client_order_id);
2004 } else {
2005 let mapped_client_order_id = self
2006 .cache
2007 .borrow()
2008 .client_order_id(&report.venue_order_id)
2009 .copied();
2010
2011 if let Some(client_order_id) = mapped_client_order_id {
2015 venue_reported_ids.insert(client_order_id);
2016 self.missing_order_coverage_warnings
2017 .shift_remove(&client_order_id);
2018 self.recon_check_retries.shift_remove(&client_order_id);
2019 }
2020 }
2021 }
2022
2023 let mut events = Vec::new();
2024 let mut targeted_candidates = Vec::new();
2025
2026 for sourced in all_reports {
2027 let report = sourced.report;
2028 if let Some(client_order_id) = &report.client_order_id
2029 && let Some(order) = self.get_order(*client_order_id)
2030 {
2031 let threshold = Duration::from_nanos(self.config.open_check_threshold_ns);
2033 if let Some(elapsed) = self.order_local_activity.elapsed(client_order_id)
2034 && elapsed < threshold
2035 {
2036 let elapsed_ms = elapsed.as_millis();
2037 let threshold_ms = threshold.as_millis();
2038 log::debug!(
2039 "Deferring reconciliation for {client_order_id}: recent local activity ({elapsed_ms}ms < threshold={threshold_ms}ms)",
2040 );
2041 continue;
2042 }
2043
2044 let instrument = self.get_instrument(&report.instrument_id);
2045 let commission_client = clients
2046 .iter()
2047 .find(|client| client.client_id() == sourced.client_id)
2048 .copied();
2049
2050 match self.reconcile_order_report(
2051 &order,
2052 &report,
2053 instrument.as_ref(),
2054 commission_client,
2055 ) {
2056 Ok(Some(event)) => events.push(event),
2057 Ok(None) => {}
2058 Err(e) => log::error!(
2059 "Deferring reconciliation for {client_order_id}: venue commission calculation failed: {e}"
2060 ),
2061 }
2062 }
2063 }
2064
2065 if self.config.open_check_open_only {
2070 let cached_ids: IndexSet<ClientOrderId> = check
2071 .filtered_orders
2072 .iter()
2073 .map(Order::client_order_id)
2074 .collect();
2075 let missing_at_venue: IndexSet<ClientOrderId> = cached_ids
2076 .difference(&venue_reported_ids)
2077 .copied()
2078 .collect();
2079
2080 if !missing_at_venue.is_empty() {
2081 log::debug!(
2082 "{} cached open order{} not present in venue current response",
2083 missing_at_venue.len(),
2084 if missing_at_venue.len() == 1 {
2085 " is"
2086 } else {
2087 "s are"
2088 },
2089 );
2090
2091 for client_order_id in missing_at_venue {
2092 log::debug!("Cached open order missing from venue response: {client_order_id}");
2093 }
2094 }
2095 } else {
2096 let candidates: Vec<&OrderAny> = if let Some(cutoff) = check.start {
2097 check
2098 .filtered_orders
2099 .iter()
2100 .filter(|o| o.ts_last() >= cutoff)
2101 .collect()
2102 } else {
2103 check.filtered_orders.iter().collect()
2104 };
2105
2106 for order in candidates {
2107 let client_order_id = order.client_order_id();
2108 if venue_reported_ids.contains(&client_order_id) {
2109 continue;
2110 }
2111
2112 let coverage = check
2113 .client_coverage
2114 .get(&client_order_id)
2115 .unwrap_or(&ReportClientCoverage::Unresolved);
2116
2117 let ReportClientCoverage::Resolved(responsible_clients) = coverage else {
2118 if self.missing_order_coverage_warnings.insert(client_order_id) {
2119 log::warn!(
2120 "Skipping order reconciliation for {client_order_id}: responsible execution client coverage is unresolved"
2121 );
2122 }
2123 continue;
2124 };
2125
2126 if responsible_clients.is_empty() {
2127 if self.missing_order_coverage_warnings.insert(client_order_id) {
2128 log::warn!(
2129 "Skipping order reconciliation for {client_order_id}: responsible execution client coverage is unresolved"
2130 );
2131 }
2132 continue;
2133 }
2134
2135 let missing_clients = responsible_clients
2136 .difference(queried_clients)
2137 .copied()
2138 .collect::<IndexSet<_>>();
2139
2140 if !missing_clients.is_empty() {
2141 if self.missing_order_coverage_warnings.insert(client_order_id) {
2142 log::warn!(
2143 "Skipping order reconciliation for {client_order_id}: responsible execution clients were not queried: {missing_clients:?}"
2144 );
2145 }
2146 continue;
2147 }
2148
2149 let failed_responsible_clients = responsible_clients
2150 .intersection(failed_clients)
2151 .copied()
2152 .collect::<IndexSet<_>>();
2153
2154 if !failed_responsible_clients.is_empty() {
2155 log::warn!(
2156 "Skipping order reconciliation for {client_order_id}: failed to query responsible execution clients: {failed_responsible_clients:?}"
2157 );
2158 continue;
2159 }
2160
2161 self.missing_order_coverage_warnings
2162 .shift_remove(&client_order_id);
2163 if let Some(order) = self.prepare_missing_order_query(client_order_id) {
2164 targeted_candidates.push((order, responsible_clients.clone()));
2165 }
2166 }
2167 }
2168
2169 targeted_candidates.sort_by_key(|(order, _)| {
2170 let client_order_id = order.client_order_id();
2171 (
2172 self.order_query_recency.last_marked(&client_order_id),
2173 client_order_id,
2174 )
2175 });
2176
2177 let query_limit = self.config.max_single_order_queries_per_cycle as usize;
2178 let mut planned_queries = 0usize;
2179 let mut cap_deferred_orders = 0usize;
2180 let mut targeted_queries = Vec::new();
2181
2182 for (order, responsible_clients) in targeted_candidates {
2183 let client_order_id = order.client_order_id();
2184
2185 let required_queries = responsible_clients.len();
2186 let exceeds_query_limit = planned_queries + required_queries > query_limit;
2187 let can_run_oversized_group = planned_queries == 0 && query_limit > 0;
2188 if required_queries == 0 || (exceeds_query_limit && !can_run_oversized_group) {
2189 cap_deferred_orders += 1;
2190 continue;
2191 }
2192
2193 if required_queries > query_limit {
2194 log::warn!(
2195 "Targeted order query for {client_order_id} requires {required_queries} responsible clients, exceeding the per-cycle limit {query_limit} to avoid indefinite deferral"
2196 );
2197 }
2198
2199 planned_queries += required_queries;
2200 self.order_query_recency.mark(client_order_id);
2201 self.targeted_order_queries.insert(client_order_id);
2202 targeted_queries.push(TargetedOrderQuery {
2203 client_order_id,
2204 responsible_clients,
2205 command: GenerateOrderStatusReport::new(
2206 UUID4::new(),
2207 self.clock.borrow().timestamp_ns(),
2208 Some(order.instrument_id()),
2209 Some(client_order_id),
2210 order.venue_order_id(),
2211 None,
2212 None,
2213 ),
2214 });
2215 }
2216
2217 if cap_deferred_orders > 0 {
2218 log::warn!(
2219 "Reached max single-order queries ({query_limit}) this cycle, deferring {cap_deferred_orders} order(s)"
2220 );
2221 }
2222
2223 OpenOrderReconciliationResult {
2224 events,
2225 targeted_queries,
2226 }
2227 }
2228
2229 pub(crate) fn reconcile_targeted_order_reports(
2230 &mut self,
2231 results: Vec<TargetedOrderReportResult>,
2232 clients: &[&dyn ExecutionClient],
2233 ) -> Vec<OrderEventAny> {
2234 let mut events = Vec::new();
2235
2236 for result in results {
2237 let client_order_id = result.client_order_id;
2238 self.targeted_order_queries.shift_remove(&client_order_id);
2239
2240 if let Some(report) = result.report {
2241 self.recon_check_retries.shift_remove(&client_order_id);
2242 self.missing_order_coverage_warnings
2243 .shift_remove(&client_order_id);
2244
2245 let Some(order) = self.get_order(client_order_id) else {
2246 continue;
2247 };
2248 let instrument = self.get_instrument(&report.instrument_id);
2249 let commission_client = result.client_id.and_then(|client_id| {
2250 clients
2251 .iter()
2252 .find(|client| client.client_id() == client_id)
2253 .copied()
2254 });
2255
2256 log::info!(
2257 color = LogColor::Blue as u8;
2258 "Found {client_order_id} via targeted order status query: {}",
2259 report.order_status,
2260 );
2261
2262 match self.reconcile_order_report(
2263 &order,
2264 &report,
2265 instrument.as_ref(),
2266 commission_client,
2267 ) {
2268 Ok(Some(event)) => events.push(event),
2269 Ok(None) => {}
2270 Err(e) => log::error!(
2271 "Deferring targeted reconciliation for {client_order_id}: venue commission calculation failed: {e}"
2272 ),
2273 }
2274 continue;
2275 }
2276
2277 if result.coverage_complete {
2278 events.extend(self.resolve_missing_order(client_order_id));
2279 } else {
2280 log::warn!(
2281 "Deferring missing-order resolution for {client_order_id}: targeted order status coverage was incomplete"
2282 );
2283 }
2284 }
2285
2286 events
2287 }
2288
2289 #[must_use]
2291 pub(crate) fn prepare_position_report_check(
2292 &self,
2293 command_id: UUID4,
2294 clients: &[&dyn ExecutionClient],
2295 ) -> PositionReportCheck {
2296 let position_keys = self.open_position_keys_for_reconciliation();
2297 let client_coverage = position_keys
2298 .iter()
2299 .map(|key| {
2300 (
2301 *key,
2302 Self::resolve_position_report_client_coverage(*key, clients),
2303 )
2304 })
2305 .collect();
2306 let activity_revisions = position_keys
2307 .iter()
2308 .map(|key| (*key, self.position_activity_revision(key)))
2309 .collect();
2310
2311 log::debug!(
2312 "Found {} unique instrument/account combination{} with open positions",
2313 position_keys.len(),
2314 if position_keys.len() == 1 { "" } else { "s" }
2315 );
2316
2317 let mut command = GeneratePositionStatusReports::new(
2318 command_id,
2319 self.clock.borrow().timestamp_ns(),
2320 None, None, None, None, None, );
2326 command.log_receipt_level = LogLevel::Debug;
2327
2328 PositionReportCheck {
2329 command,
2330 client_coverage,
2331 activity_revisions,
2332 }
2333 }
2334
2335 fn resolve_position_report_client_coverage(
2336 key: InstrumentAccountKey,
2337 clients: &[&dyn ExecutionClient],
2338 ) -> ReportClientCoverage {
2339 let account_clients = clients
2340 .iter()
2341 .filter(|client| client.account_id() == key.1)
2342 .map(|client| client.client_id())
2343 .collect::<IndexSet<_>>();
2344
2345 if !account_clients.is_empty() {
2346 return ReportClientCoverage::Resolved(account_clients);
2347 }
2348
2349 let venue_clients = clients
2350 .iter()
2351 .filter(|client| client.handles_order_venue(key.0.venue))
2352 .map(|client| client.client_id())
2353 .collect::<IndexSet<_>>();
2354
2355 if venue_clients.is_empty() {
2356 ReportClientCoverage::Unresolved
2357 } else {
2358 ReportClientCoverage::Resolved(venue_clients)
2359 }
2360 }
2361
2362 pub async fn check_positions_consistency(
2372 &mut self,
2373 clients: &[&dyn ExecutionClient],
2374 ) -> Vec<OrderEventAny> {
2375 let check = self.prepare_position_report_check(UUID4::new(), clients);
2376 let mut reports = Vec::new();
2377 let mut queried_clients = IndexSet::new();
2378 let mut failed_clients = IndexSet::new();
2379
2380 for client in clients {
2381 let client_id = client.client_id();
2382 queried_clients.insert(client_id);
2383 self.set_position_reconciliation_tolerance(
2384 client.account_id(),
2385 client.position_reconciliation_tolerance(),
2386 );
2387
2388 match client
2389 .generate_position_status_reports(&check.command)
2390 .await
2391 {
2392 Ok(client_reports) => {
2393 reports.extend(client_reports);
2394 }
2395 Err(e) => {
2396 failed_clients.insert(client_id);
2397 log::warn!(
2398 "Failed to query position reports from {}: {e}",
2399 client.client_id()
2400 );
2401 }
2402 }
2403 }
2404
2405 self.reconcile_position_reports(&check, reports, &queried_clients, &failed_clients)
2406 }
2407
2408 #[must_use]
2410 pub(crate) fn reconcile_position_reports(
2411 &mut self,
2412 check: &PositionReportCheck,
2413 reports: Vec<PositionStatusReport>,
2414 queried_clients: &IndexSet<ClientId>,
2415 failed_clients: &IndexSet<ClientId>,
2416 ) -> Vec<OrderEventAny> {
2417 log::debug!("Checking position consistency between cached-state and venues");
2418
2419 let mut venue_positions: IndexMap<InstrumentAccountKey, Vec<PositionStatusReport>> =
2420 IndexMap::new();
2421
2422 for report in reports {
2423 if !self.should_reconcile_instrument(&report.instrument_id) {
2424 continue;
2425 }
2426
2427 venue_positions
2428 .entry((report.instrument_id, report.account_id))
2429 .or_default()
2430 .push(report);
2431 }
2432
2433 let mut events = Vec::new();
2434
2435 for key in check.client_coverage.keys() {
2436 let prepared_revision = check
2437 .activity_revisions
2438 .get(key)
2439 .copied()
2440 .unwrap_or_default();
2441
2442 if self.position_activity_revision(key) > prepared_revision {
2443 log::debug!(
2444 "Deferring position reconciliation for {}/{}: local activity recorded during report request",
2445 key.0,
2446 key.1,
2447 );
2448 continue;
2449 }
2450
2451 let venue_reports = venue_positions
2452 .get(key)
2453 .map(Vec::as_slice)
2454 .unwrap_or_default();
2455
2456 if venue_reports.is_empty() {
2457 match check.client_coverage.get(key) {
2458 Some(ReportClientCoverage::Resolved(responsible_clients))
2459 if !responsible_clients.is_empty()
2460 && responsible_clients.is_subset(queried_clients)
2461 && responsible_clients.is_disjoint(failed_clients) => {}
2462 Some(ReportClientCoverage::Resolved(responsible_clients))
2463 if responsible_clients.is_empty() =>
2464 {
2465 log::warn!(
2466 "Skipping position reconciliation for {}/{}: responsible execution client coverage is unresolved",
2467 key.0,
2468 key.1,
2469 );
2470 continue;
2471 }
2472 Some(ReportClientCoverage::Resolved(responsible_clients))
2473 if !responsible_clients.is_subset(queried_clients) =>
2474 {
2475 log::warn!(
2476 "Skipping position reconciliation for {}/{}: responsible execution clients were not all queried",
2477 key.0,
2478 key.1,
2479 );
2480 continue;
2481 }
2482 Some(ReportClientCoverage::Resolved(responsible_clients)) => {
2483 let failed_responsible_clients = responsible_clients
2484 .intersection(failed_clients)
2485 .copied()
2486 .collect::<IndexSet<_>>();
2487 log::warn!(
2488 "Skipping position reconciliation for {}/{}: failed to query responsible execution clients: {failed_responsible_clients:?}",
2489 key.0,
2490 key.1,
2491 );
2492 continue;
2493 }
2494 Some(ReportClientCoverage::Unresolved) | None => {
2495 log::warn!(
2496 "Skipping position reconciliation for {}/{}: responsible execution client coverage is unresolved",
2497 key.0,
2498 key.1,
2499 );
2500 continue;
2501 }
2502 }
2503 }
2504
2505 if let Some(discrepancy_events) = self.check_position_discrepancy(*key, venue_reports) {
2506 events.extend(discrepancy_events);
2507 }
2508 }
2509
2510 let current_position_keys = self.open_position_keys_for_reconciliation();
2511
2512 for (key, venue_reports) in &venue_positions {
2513 if check.client_coverage.contains_key(key)
2514 || venue_reports
2515 .iter()
2516 .all(|report| report.signed_decimal_qty == Decimal::ZERO)
2517 {
2518 continue;
2519 }
2520
2521 if current_position_keys.contains(key) {
2522 log::debug!(
2523 "Deferring position reconciliation for {}/{}: position opened after client coverage was recorded",
2524 key.0,
2525 key.1,
2526 );
2527 continue;
2528 }
2529
2530 if let Some(discrepancy_events) = self.check_position_discrepancy(*key, venue_reports) {
2531 events.extend(discrepancy_events);
2532 }
2533 }
2534
2535 let active_keys: IndexSet<InstrumentAccountKey> = current_position_keys
2538 .into_iter()
2539 .chain(
2540 venue_positions
2541 .iter()
2542 .filter(|(_, reports)| {
2543 reports
2544 .iter()
2545 .any(|report| report.signed_decimal_qty != Decimal::ZERO)
2546 })
2547 .map(|(k, _)| *k),
2548 )
2549 .collect();
2550 self.position_reconciliation_states
2551 .retain(|k, _| active_keys.contains(k));
2552
2553 events
2554 }
2555
2556 fn positions_avg_px(cached_positions: &[Position]) -> Option<Decimal> {
2557 let mut total_value = Decimal::ZERO;
2558 let mut total_qty = Decimal::ZERO;
2559
2560 for position in cached_positions {
2561 let qty = position.signed_decimal_qty().abs();
2562 if position.avg_px_open > 0.0
2563 && qty > Decimal::ZERO
2564 && let Ok(avg_px) = Decimal::from_str(&position.avg_px_open.to_string())
2565 {
2566 total_value += avg_px * qty;
2567 total_qty += qty;
2568 }
2569 }
2570
2571 if total_qty > Decimal::ZERO {
2572 Some(total_value / total_qty)
2573 } else {
2574 None
2575 }
2576 }
2577
2578 pub fn register_inflight(&mut self, client_order_id: ClientOrderId) {
2580 if self
2581 .config
2582 .filtered_client_order_ids
2583 .contains(&client_order_id)
2584 {
2585 return;
2586 }
2587
2588 self.inflight_checks.insert(
2589 client_order_id,
2590 InflightCheck {
2591 client_order_id,
2592 submitted_at: dst::time::Instant::now(),
2593 retry_count: 0,
2594 last_query_at: None,
2595 },
2596 );
2597 self.recon_check_retries.insert(client_order_id, 0);
2598 self.order_query_recency.remove(&client_order_id);
2599 self.order_local_activity.remove(&client_order_id);
2600 }
2601
2602 pub fn record_local_activity(&mut self, client_order_id: ClientOrderId) {
2609 self.order_local_activity.mark(client_order_id);
2610 }
2611
2612 pub fn clear_recon_tracking(&mut self, client_order_id: &ClientOrderId, drop_last_query: bool) {
2614 self.inflight_checks.shift_remove(client_order_id);
2615 self.recon_check_retries.shift_remove(client_order_id);
2616 self.missing_order_coverage_warnings
2617 .shift_remove(client_order_id);
2618 self.unresolved_order_coverage.shift_remove(client_order_id);
2619 self.targeted_order_queries.shift_remove(client_order_id);
2620
2621 if drop_last_query {
2622 self.order_query_recency.remove(client_order_id);
2623 }
2624 self.order_local_activity.remove(client_order_id);
2625 }
2626
2627 #[cfg(feature = "node")]
2628 pub(crate) fn remove_targeted_order_queries(&mut self, client_order_ids: &[ClientOrderId]) {
2629 for client_order_id in client_order_ids {
2630 self.targeted_order_queries.shift_remove(client_order_id);
2631 }
2632 }
2633
2634 #[must_use]
2636 pub fn get_external_order_claim(&self, instrument_id: &InstrumentId) -> Option<StrategyId> {
2637 self.external_order_claims.get(instrument_id).copied()
2638 }
2639
2640 #[must_use]
2642 #[cfg(feature = "node")]
2643 pub(crate) fn get_external_order_claims_for_strategy(
2644 &self,
2645 strategy_id: StrategyId,
2646 ) -> HashSet<InstrumentId> {
2647 self.external_order_claims
2648 .iter()
2649 .filter_map(|(instrument_id, owner)| (*owner == strategy_id).then_some(*instrument_id))
2650 .collect()
2651 }
2652
2653 pub fn claim_external_orders(
2659 &mut self,
2660 instrument_id: InstrumentId,
2661 strategy_id: StrategyId,
2662 ) -> anyhow::Result<()> {
2663 if let Some(existing) = self.external_order_claims.get(&instrument_id) {
2664 anyhow::bail!("External order claim for {instrument_id} already exists for {existing}");
2665 }
2666
2667 self.external_order_claims
2668 .insert(instrument_id, strategy_id);
2669 Ok(())
2670 }
2671
2672 #[cfg(feature = "node")]
2673 pub(crate) fn register_external_order_claims(
2674 &mut self,
2675 strategy_id: StrategyId,
2676 instrument_ids: &HashSet<InstrumentId>,
2677 ) {
2678 self.external_order_claims.extend(
2679 instrument_ids
2680 .iter()
2681 .map(|instrument_id| (*instrument_id, strategy_id)),
2682 );
2683 }
2684
2685 #[cfg(feature = "node")]
2691 pub(crate) fn deregister_external_order_claims(&mut self, strategy_id: StrategyId) {
2692 self.external_order_claims
2693 .retain(|_, owner| *owner != strategy_id);
2694 }
2695
2696 pub fn record_position_activity(&mut self, instrument_id: InstrumentId, account_id: AccountId) {
2710 let key = (instrument_id, account_id);
2711 self.position_local_activity.mark(key);
2712 let revision = self
2713 .position_local_activity_revisions
2714 .entry(key)
2715 .or_default();
2716 *revision = revision.saturating_add(1);
2717 }
2718
2719 fn position_activity_revision(&self, key: &InstrumentAccountKey) -> u64 {
2720 self.position_local_activity_revisions
2721 .get(key)
2722 .copied()
2723 .unwrap_or_default()
2724 }
2725
2726 #[must_use]
2729 pub fn position_recon_retry_count(&self, key: &InstrumentAccountKey) -> u32 {
2730 self.position_reconciliation_states
2731 .get(key)
2732 .map_or(0, |state| state.retries)
2733 }
2734
2735 #[must_use]
2738 pub fn recon_check_retry_count(&self, client_order_id: &ClientOrderId) -> u32 {
2739 self.recon_check_retries
2740 .get(client_order_id)
2741 .copied()
2742 .unwrap_or(0)
2743 }
2744
2745 pub fn observe_order_event(&mut self, event: &OrderEventAny) {
2755 match event {
2756 OrderEventAny::Filled(fill) => {
2757 self.record_position_activity(fill.instrument_id, fill.account_id);
2758 }
2759 OrderEventAny::Accepted(_)
2760 | OrderEventAny::Rejected(_)
2761 | OrderEventAny::Canceled(_)
2762 | OrderEventAny::Expired(_)
2763 | OrderEventAny::Denied(_)
2764 | OrderEventAny::Updated(_)
2765 | OrderEventAny::ModifyRejected(_)
2766 | OrderEventAny::CancelRejected(_) => {
2767 self.clear_recon_tracking(&event.client_order_id(), true);
2768 }
2769 _ => {}
2770 }
2771
2772 self.record_local_activity(event.client_order_id());
2773 }
2774
2775 pub fn observe_execution_report(&mut self, report: &ExecutionReport) {
2787 match report {
2788 ExecutionReport::Order(order_report) => {
2789 self.observe_order_status_report(order_report);
2790 }
2791 ExecutionReport::Fill(fill_report) => {
2792 let client_order_id = fill_report.client_order_id.or_else(|| {
2793 self.cache
2794 .borrow()
2795 .client_order_id(&fill_report.venue_order_id)
2796 .copied()
2797 });
2798
2799 if let Some(coid) = client_order_id {
2800 self.record_local_activity(coid);
2801 }
2802 self.record_position_activity(fill_report.instrument_id, fill_report.account_id);
2803 }
2804 ExecutionReport::OrderWithFills(order_report, fills) => {
2805 self.observe_order_status_report(order_report);
2806
2807 for fill_report in fills {
2808 self.record_position_activity(
2809 fill_report.instrument_id,
2810 fill_report.account_id,
2811 );
2812 }
2813 }
2814 ExecutionReport::Position(position_report) => {
2815 self.record_position_activity(
2816 position_report.instrument_id,
2817 position_report.account_id,
2818 );
2819 }
2820 ExecutionReport::MassStatus(_) => {
2821 }
2823 }
2824 }
2825
2826 fn observe_order_status_report(&mut self, report: &OrderStatusReport) {
2827 let Some(client_order_id) = report.client_order_id else {
2828 return;
2829 };
2830
2831 if !matches!(
2832 report.order_status,
2833 OrderStatus::PendingUpdate | OrderStatus::PendingCancel
2834 ) {
2835 self.clear_recon_tracking(&client_order_id, report.order_status.is_closed());
2836 }
2837
2838 self.record_local_activity(client_order_id);
2842 }
2843
2844 #[must_use]
2846 pub fn is_fill_recently_processed(
2847 &self,
2848 account_id: AccountId,
2849 instrument_id: InstrumentId,
2850 trade_id: TradeId,
2851 ) -> bool {
2852 self.recent_fills_cache
2853 .contains_key(&(account_id, instrument_id, trade_id))
2854 }
2855
2856 pub fn mark_fill_processed(
2858 &mut self,
2859 account_id: AccountId,
2860 instrument_id: InstrumentId,
2861 trade_id: TradeId,
2862 ) {
2863 self.recent_fills_cache
2864 .mark((account_id, instrument_id, trade_id));
2865 }
2866
2867 pub fn commit_recent_fill_if_applied(&mut self, fill: &OrderFilled) {
2869 let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
2870 if self.is_fill_applied(fill, fill_key) {
2871 self.mark_fill_processed(fill_key.0, fill_key.1, fill_key.2);
2872 }
2873 }
2874
2875 pub fn prune_recent_fills_cache(&mut self, ttl_secs: f64) {
2879 let ttl = match Duration::try_from_secs_f64(ttl_secs) {
2887 Ok(ttl) => ttl,
2888 Err(_) if ttl_secs > 0.0 => Duration::MAX,
2889 Err(_) => Duration::ZERO,
2890 };
2891
2892 self.recent_fills_cache.prune_older_than(ttl);
2893 }
2894
2895 pub fn prune_processed_fills(&mut self) {
2900 let Some(lookback_mins) = self.config.lookback_mins else {
2901 return;
2902 };
2903
2904 let ttl = Duration::from_mins(lookback_mins).max(Duration::from_mins(1));
2905 self.processed_fills.prune_older_than(ttl);
2906 }
2907
2908 pub fn prune_order_local_activity(&mut self) {
2910 self.order_local_activity
2911 .prune_older_than(Duration::from_nanos(self.config.open_check_threshold_ns));
2912 }
2913
2914 pub fn purge_closed_orders(&mut self) {
2916 let Some(buffer_mins) = self.config.purge_closed_orders_buffer_mins else {
2917 return;
2918 };
2919
2920 let ts_now = self.clock.borrow().timestamp_ns();
2921 let buffer_secs = mins_to_secs(u64::from(buffer_mins));
2922
2923 self.cache
2924 .borrow_mut()
2925 .purge_closed_orders(ts_now, buffer_secs);
2926 }
2927
2928 pub fn purge_closed_positions(&mut self) {
2930 let Some(buffer_mins) = self.config.purge_closed_positions_buffer_mins else {
2931 return;
2932 };
2933
2934 let ts_now = self.clock.borrow().timestamp_ns();
2935 let buffer_secs = mins_to_secs(u64::from(buffer_mins));
2936
2937 self.cache
2938 .borrow_mut()
2939 .purge_closed_positions(ts_now, buffer_secs);
2940 }
2941
2942 pub fn purge_account_events(&mut self) {
2944 let Some(lookback_mins) = self.config.purge_account_events_lookback_mins else {
2945 return;
2946 };
2947
2948 let ts_now = self.clock.borrow().timestamp_ns();
2949 let lookback_secs = mins_to_secs(u64::from(lookback_mins));
2950
2951 self.cache
2952 .borrow_mut()
2953 .purge_account_events(ts_now, lookback_secs);
2954 }
2955
2956 fn get_order(&self, client_order_id: ClientOrderId) -> Option<OrderAny> {
2959 self.cache
2960 .borrow()
2961 .order(&client_order_id)
2962 .map(|o| o.clone())
2963 }
2964
2965 fn get_order_by_venue_order_id(&self, venue_order_id: VenueOrderId) -> Option<OrderAny> {
2966 let cache = self.cache.borrow();
2967 cache
2968 .client_order_id(&venue_order_id)
2969 .and_then(|client_order_id| cache.order(client_order_id).map(|o| o.clone()))
2970 }
2971
2972 fn get_instrument(&self, instrument_id: &InstrumentId) -> Option<InstrumentAny> {
2973 self.cache.borrow().instrument(instrument_id).cloned()
2974 }
2975
2976 fn should_skip_order_report(&self, report: &OrderStatusReport) -> bool {
2977 if let Some(client_order_id) = &report.client_order_id
2978 && self
2979 .config
2980 .filtered_client_order_ids
2981 .contains(client_order_id)
2982 {
2983 log::debug!(
2984 "Skipping order report {client_order_id}: in filtered_client_order_ids list"
2985 );
2986 return true;
2987 }
2988
2989 if !self.should_reconcile_instrument(&report.instrument_id) {
2990 log::debug!(
2991 "Skipping order report for {}: not in reconciliation_instrument_ids",
2992 report.instrument_id
2993 );
2994 return true;
2995 }
2996
2997 false
2998 }
2999
3000 fn should_reconcile_instrument(&self, instrument_id: &InstrumentId) -> bool {
3001 self.config.reconciliation_instrument_ids.is_empty()
3002 || self
3003 .config
3004 .reconciliation_instrument_ids
3005 .contains(instrument_id)
3006 }
3007
3008 fn prepare_missing_order_query(&mut self, client_order_id: ClientOrderId) -> Option<OrderAny> {
3009 let order = self.get_order(client_order_id)?;
3010
3011 if order.status().is_closed() {
3015 log::debug!(
3016 "Skipping missing-order resolution for {client_order_id}: already {}",
3017 order.status()
3018 );
3019 self.clear_recon_tracking(&client_order_id, true);
3020 return None;
3021 }
3022
3023 if self.order_local_activity.within(
3027 &client_order_id,
3028 Duration::from_nanos(self.config.open_check_threshold_ns),
3029 ) {
3030 return None;
3031 }
3032
3033 let retries = self.recon_check_retries.entry(client_order_id).or_insert(0);
3034 *retries = retries.saturating_add(1);
3035
3036 if *retries < self.config.open_check_missing_retries {
3037 log::debug!(
3038 "Order {} not found at venue, retry {}/{}",
3039 client_order_id,
3040 retries,
3041 self.config.open_check_missing_retries
3042 );
3043 return None;
3044 }
3045
3046 Some(order)
3047 }
3048
3049 fn resolve_missing_order(&mut self, client_order_id: ClientOrderId) -> Vec<OrderEventAny> {
3050 let mut events = Vec::new();
3051
3052 let Some(order) = self.get_order(client_order_id) else {
3053 return events;
3054 };
3055
3056 if order.status().is_closed() {
3057 log::debug!(
3058 "Skipping missing-order resolution for {client_order_id}: already {}",
3059 order.status()
3060 );
3061 self.clear_recon_tracking(&client_order_id, true);
3062 return events;
3063 }
3064
3065 if self.order_local_activity.within(
3066 &client_order_id,
3067 Duration::from_nanos(self.config.open_check_threshold_ns),
3068 ) {
3069 log::debug!(
3070 "Deferring missing-order resolution for {client_order_id}: recent local activity"
3071 );
3072 return events;
3073 }
3074
3075 let retries = self
3076 .recon_check_retries
3077 .get(&client_order_id)
3078 .copied()
3079 .unwrap_or_default();
3080 let ts_now = self.clock.borrow().timestamp_ns();
3081
3082 match order.status() {
3083 OrderStatus::Accepted | OrderStatus::Submitted => {
3084 log::warn!(
3085 "Order {client_order_id} not found at venue after {retries} retries and a targeted query, marking as REJECTED"
3086 );
3087
3088 if let Some(rejected) =
3089 create_reconciliation_rejected(&order, Some("NOT_FOUND_AT_VENUE"), ts_now)
3090 {
3091 events.push(rejected);
3092 }
3093 }
3094 OrderStatus::PartiallyFilled => {
3095 log::warn!(
3096 "Order {client_order_id} not found at venue after {retries} retries and a targeted query, marking as CANCELED"
3097 );
3098 events.push(OrderEventAny::Canceled(OrderCanceled::new(
3099 order.trader_id(),
3100 order.strategy_id(),
3101 order.instrument_id(),
3102 client_order_id,
3103 UUID4::new(),
3104 ts_now,
3105 ts_now,
3106 true,
3107 order.venue_order_id(),
3108 order.account_id(),
3109 )));
3110 }
3111 OrderStatus::PendingUpdate | OrderStatus::PendingCancel => {
3112 log::debug!(
3113 "Deferring resolution for {client_order_id}: still inflight as {}",
3114 order.status()
3115 );
3116 self.recon_check_retries.shift_remove(&client_order_id);
3125 if let Some(check) = self.inflight_checks.get_mut(&client_order_id) {
3126 check.retry_count = 0;
3127 check.last_query_at = Some(dst::time::Instant::now());
3128 }
3129 self.order_query_recency.mark(client_order_id);
3130 return events;
3131 }
3132 status => {
3133 log::warn!(
3134 "Skipping missing-order resolution for {client_order_id}: unexpected status {status}"
3135 );
3136 }
3137 }
3138
3139 self.clear_recon_tracking(&client_order_id, true);
3140 events
3141 }
3142
3143 fn check_position_discrepancy(
3144 &mut self,
3145 key: InstrumentAccountKey,
3146 venue_reports: &[PositionStatusReport],
3147 ) -> Option<Vec<OrderEventAny>> {
3148 let (instrument_id, account_id) = key;
3149
3150 let cached_positions = {
3151 let cache = self.cache.borrow();
3152 cache
3153 .positions_open(None, Some(&instrument_id), None, Some(&account_id), None)
3154 .into_iter()
3155 .map(|position| (*position).clone())
3156 .collect::<Vec<_>>()
3157 };
3158 let (cached_signed_qty, cached_long_qty, cached_short_qty) = Self::position_qty_aggregates(
3159 cached_positions.iter().map(Position::signed_decimal_qty),
3160 );
3161 let (venue_signed_qty, venue_long_qty, venue_short_qty) = Self::position_qty_aggregates(
3162 venue_reports.iter().map(|report| report.signed_decimal_qty),
3163 );
3164 let nonflat_count = venue_reports
3165 .iter()
3166 .filter(|report| report.signed_decimal_qty != Decimal::ZERO)
3167 .count();
3168 let venue_report = venue_reports
3169 .iter()
3170 .find(|report| report.signed_decimal_qty != Decimal::ZERO)
3171 .or_else(|| venue_reports.last());
3172
3173 let tolerance = self.position_reconciliation_tolerance(account_id);
3174 let venue_has_side_reports = venue_reports.iter().any(PositionStatusReport::is_long)
3175 && venue_reports.iter().any(PositionStatusReport::is_short);
3176 let net_qty_matches = (cached_signed_qty - venue_signed_qty).abs() <= tolerance;
3177 let side_qty_matches = (cached_long_qty - venue_long_qty).abs() <= tolerance
3178 && (cached_short_qty - venue_short_qty).abs() <= tolerance;
3179
3180 if net_qty_matches && (!venue_has_side_reports || side_qty_matches) {
3181 self.position_reconciliation_states.shift_remove(&key);
3182 return None;
3183 }
3184
3185 let ts_now = self.clock.borrow().timestamp_ns();
3186
3187 if self.position_local_activity.within(
3189 &key,
3190 Duration::from_nanos(self.config.position_check_threshold_ns),
3191 ) {
3192 log::debug!(
3193 "Skipping position reconciliation for {instrument_id}: recent activity within threshold"
3194 );
3195 return None;
3196 }
3197
3198 let report_shape = if nonflat_count > 1 || venue_has_side_reports {
3199 PositionReportShape::MultiLeg
3200 } else {
3201 PositionReportShape::Unambiguous
3202 };
3203 let retries = self
3204 .position_reconciliation_states
3205 .get(&key)
3206 .filter(|state| state.report_shape == report_shape)
3207 .map_or(0, |state| state.retries);
3208
3209 if retries >= self.config.position_check_retries {
3210 return None;
3211 }
3212
3213 if report_shape == PositionReportShape::MultiLeg {
3214 let new_retries = retries + 1;
3215 self.set_position_reconciliation_retries(key, report_shape, new_retries);
3216 log::warn!(
3217 "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})"
3218 );
3219
3220 if new_retries >= self.config.position_check_retries {
3221 log::error!(
3222 "Position discrepancy for {instrument_id}/{account_id} unresolved after {} attempts; no further reconciliation attempts will be made for the current report shape",
3223 self.config.position_check_retries,
3224 );
3225 }
3226 return None;
3227 }
3228
3229 log::warn!(
3230 "Position discrepancy detected for {instrument_id}: cached_signed_qty={cached_signed_qty}, venue_signed_qty={venue_signed_qty}"
3231 );
3232
3233 let Some(instrument) = self.cache.borrow().instrument(&instrument_id).cloned() else {
3234 log::debug!("Cannot reconcile position for {instrument_id}: instrument not in cache");
3235 let new_retries = retries + 1;
3236 self.set_position_reconciliation_retries(key, report_shape, new_retries);
3237 if new_retries >= self.config.position_check_retries {
3238 log::error!(
3239 "Position discrepancy for {instrument_id} unresolved after {} attempts \
3240 (cached_qty={cached_signed_qty}, venue_qty={venue_signed_qty}); \
3241 no further reconciliation attempts will be made for the current report shape",
3242 self.config.position_check_retries,
3243 );
3244 }
3245 return None;
3246 };
3247
3248 let cached_avg_px = Self::positions_avg_px(&cached_positions);
3249 let venue_avg_px = venue_report.and_then(|r| r.avg_px_open);
3250
3251 let crosses_zero = (cached_signed_qty > Decimal::ZERO && venue_signed_qty < Decimal::ZERO)
3252 || (cached_signed_qty < Decimal::ZERO && venue_signed_qty > Decimal::ZERO);
3253
3254 let result = if crosses_zero {
3255 let venue_ts_last = venue_report.map_or(ts_now, |r| r.ts_last);
3256 self.reconcile_cross_zero_position(
3257 &instrument,
3258 account_id,
3259 instrument_id,
3260 cached_signed_qty,
3261 cached_avg_px,
3262 venue_signed_qty,
3263 venue_avg_px,
3264 ts_now,
3265 venue_ts_last,
3266 )
3267 } else {
3268 let qty_diff = venue_signed_qty - cached_signed_qty;
3269 let order_side = if qty_diff > Decimal::ZERO {
3270 OrderSide::Buy
3271 } else {
3272 OrderSide::Sell
3273 };
3274
3275 let reconciliation_px = calculate_reconciliation_price(
3276 cached_signed_qty,
3277 cached_avg_px,
3278 venue_signed_qty,
3279 venue_avg_px,
3280 );
3281
3282 match reconciliation_px.or(venue_avg_px).or(cached_avg_px) {
3283 Some(fill_px) => {
3284 let fill_qty = qty_diff.abs();
3285 let venue_position_id = venue_report.and_then(|r| r.venue_position_id);
3286 let venue_ts_last = venue_report.map_or(ts_now, |r| r.ts_last);
3287
3288 Quantity::from_decimal_dp(fill_qty, instrument.size_precision())
3289 .ok()
3290 .map(|order_qty| {
3291 let fill_price =
3292 Price::from_decimal_dp(fill_px, instrument.price_precision()).ok();
3293 let venue_order_id = create_position_reconciliation_venue_order_id(
3294 account_id,
3295 instrument_id,
3296 order_side,
3297 OrderType::Market,
3298 order_qty,
3299 fill_price,
3300 venue_position_id,
3301 None,
3302 venue_ts_last,
3303 );
3304
3305 OrderStatusReport::new(
3306 account_id,
3307 instrument_id,
3308 None,
3309 venue_order_id,
3310 order_side,
3311 OrderType::Market,
3312 TimeInForce::Gtc,
3313 OrderStatus::Filled,
3314 order_qty,
3315 order_qty,
3316 ts_now,
3317 ts_now,
3318 ts_now,
3319 None,
3320 )
3321 .with_avg_px(fill_px)
3322 })
3323 .map(|order_report| {
3324 log::info!(
3325 color = LogColor::Blue as u8;
3326 "Generating synthetic fill for position reconciliation {instrument_id}: side={order_side:?}, qty={}, px={fill_px}", qty_diff.abs(),
3327 );
3328
3329 let (events, _) = self.handle_external_order(
3330 &order_report,
3331 account_id,
3332 &instrument,
3333 &[],
3334 true,
3335 None,
3336 None,
3337 );
3338 events
3339 })
3340 }
3341 None => None,
3342 }
3343 };
3344
3345 if result.is_none() || result.as_ref().is_some_and(Vec::is_empty) {
3347 let new_retries = retries + 1;
3348 self.set_position_reconciliation_retries(key, report_shape, new_retries);
3349 if new_retries >= self.config.position_check_retries {
3350 log::error!(
3351 "Position discrepancy for {} unresolved after {} attempts \
3352 (cached_qty={}, venue_qty={}); \
3353 no further reconciliation attempts will be made for the current report shape",
3354 instrument_id,
3355 self.config.position_check_retries,
3356 cached_signed_qty,
3357 venue_signed_qty,
3358 );
3359 }
3360 } else {
3361 self.position_reconciliation_states.shift_remove(&key);
3362 }
3363
3364 result
3365 }
3366
3367 fn set_position_reconciliation_retries(
3368 &mut self,
3369 key: InstrumentAccountKey,
3370 report_shape: PositionReportShape,
3371 retries: u32,
3372 ) {
3373 self.position_reconciliation_states.insert(
3374 key,
3375 PositionReconciliationState {
3376 report_shape,
3377 retries,
3378 },
3379 );
3380 }
3381
3382 fn position_qty_aggregates(
3383 signed_quantities: impl Iterator<Item = Decimal>,
3384 ) -> (Decimal, Decimal, Decimal) {
3385 signed_quantities.fold(
3386 (Decimal::ZERO, Decimal::ZERO, Decimal::ZERO),
3387 |(net, long, short), qty| {
3388 if qty > Decimal::ZERO {
3389 (net + qty, long + qty, short)
3390 } else {
3391 (net + qty, long, short + qty.abs())
3392 }
3393 },
3394 )
3395 }
3396
3397 #[expect(clippy::too_many_arguments)]
3400 fn reconcile_cross_zero_position(
3401 &self,
3402 instrument: &InstrumentAny,
3403 account_id: AccountId,
3404 instrument_id: InstrumentId,
3405 cached_signed_qty: Decimal,
3406 cached_avg_px: Option<Decimal>,
3407 venue_signed_qty: Decimal,
3408 venue_avg_px: Option<Decimal>,
3409 ts_now: UnixNanos,
3410 venue_ts_last: UnixNanos,
3411 ) -> Option<Vec<OrderEventAny>> {
3412 log::info!(
3413 color = LogColor::Blue as u8;
3414 "Position crosses zero for {instrument_id}: cached={cached_signed_qty}, venue={venue_signed_qty}. Splitting into two fills",
3415 );
3416
3417 let close_qty = cached_signed_qty.abs();
3418 let close_side = if cached_signed_qty < Decimal::ZERO {
3419 OrderSide::Buy } else {
3421 OrderSide::Sell };
3423 let open_qty = venue_signed_qty.abs();
3424 let open_side = if venue_signed_qty > Decimal::ZERO {
3425 OrderSide::Buy } else {
3427 OrderSide::Sell };
3429
3430 let Some(close_px) = cached_avg_px else {
3431 log::warn!("Cannot close position for {instrument_id}: no cached average price");
3432 return None;
3433 };
3434
3435 let open_report = match venue_avg_px {
3436 Some(open_px) => Some((
3437 build_cross_zero_leg_report(
3438 instrument,
3439 account_id,
3440 instrument_id,
3441 open_side,
3442 open_qty,
3443 open_px,
3444 "OPEN",
3445 ts_now,
3446 venue_ts_last,
3447 )?,
3448 open_px,
3449 )),
3450 None => None,
3451 };
3452
3453 let close_report = build_cross_zero_leg_report(
3454 instrument,
3455 account_id,
3456 instrument_id,
3457 close_side,
3458 close_qty,
3459 close_px,
3460 "CLOSE",
3461 ts_now,
3462 venue_ts_last,
3463 )?;
3464
3465 log::info!(
3466 color = LogColor::Blue as u8;
3467 "Generating close fill for cross-zero {instrument_id}: side={close_side:?}, qty={close_qty}, px={close_px}",
3468 );
3469
3470 let (close_events, _) = self.handle_external_order(
3471 &close_report,
3472 account_id,
3473 instrument,
3474 &[],
3475 true,
3476 None,
3477 None,
3478 );
3479 let mut all_events = close_events;
3480
3481 if let Some((open_report, open_px)) = open_report {
3482 log::info!(
3483 color = LogColor::Blue as u8;
3484 "Generating open fill for cross-zero {instrument_id}: side={open_side:?}, qty={open_qty}, px={open_px}",
3485 );
3486
3487 let (open_events, _) = self.handle_external_order(
3488 &open_report,
3489 account_id,
3490 instrument,
3491 &[],
3492 true,
3493 None,
3494 None,
3495 );
3496 all_events.extend(open_events);
3497 } else {
3498 log::warn!("Cannot open new position for {instrument_id}: no venue average price");
3499 }
3500
3501 Some(all_events)
3502 }
3503
3504 fn create_position_from_report(
3509 &self,
3510 report: &PositionStatusReport,
3511 account_id: AccountId,
3512 instrument: &InstrumentAny,
3513 ) -> Option<Vec<OrderEventAny>> {
3514 let instrument_id = report.instrument_id;
3515 let venue_signed_qty = report.signed_decimal_qty;
3516
3517 if venue_signed_qty == Decimal::ZERO {
3518 return None;
3519 }
3520
3521 let order_side = if venue_signed_qty > Decimal::ZERO {
3522 OrderSide::Buy
3523 } else {
3524 OrderSide::Sell
3525 };
3526
3527 let qty_abs = venue_signed_qty.abs();
3528 let venue_avg_px = report.avg_px_open?;
3529
3530 let ts_now = self.clock.borrow().timestamp_ns();
3531 let order_qty = Quantity::from_decimal_dp(qty_abs, instrument.size_precision()).ok()?;
3532 let fill_price = Price::from_decimal_dp(venue_avg_px, instrument.price_precision()).ok();
3533 let venue_order_id = create_position_reconciliation_venue_order_id(
3534 account_id,
3535 instrument_id,
3536 order_side,
3537 OrderType::Market,
3538 order_qty,
3539 fill_price,
3540 report.venue_position_id,
3541 None,
3542 report.ts_last,
3543 );
3544
3545 let mut order_report = OrderStatusReport::new(
3546 account_id,
3547 instrument_id,
3548 None,
3549 venue_order_id,
3550 order_side,
3551 OrderType::Market,
3552 TimeInForce::Gtc,
3553 OrderStatus::Filled,
3554 order_qty,
3555 order_qty,
3556 ts_now,
3557 ts_now,
3558 ts_now,
3559 None,
3560 )
3561 .with_avg_px(venue_avg_px);
3562
3563 if let Some(venue_position_id) = report.venue_position_id {
3565 order_report = order_report.with_venue_position_id(venue_position_id);
3566 }
3567
3568 log::info!(
3569 color = LogColor::Blue as u8;
3570 "Creating position from venue report for {instrument_id}: side={order_side:?}, qty={qty_abs}, avg_px={venue_avg_px}",
3571 );
3572
3573 let (events, _) = self.handle_external_order(
3574 &order_report,
3575 account_id,
3576 instrument,
3577 &[],
3578 true,
3579 None,
3580 None,
3581 );
3582 Some(events)
3583 }
3584
3585 fn reconcile_position_report(
3586 &self,
3587 report: &PositionStatusReport,
3588 account_id: AccountId,
3589 instruments_with_unattributed_fills: &IndexSet<InstrumentId>,
3590 positions_with_fills: &IndexSet<PositionId>,
3591 ) -> Option<Vec<OrderEventAny>> {
3592 if report.venue_position_id.is_some() {
3593 self.reconcile_position_report_hedging(
3594 report,
3595 account_id,
3596 instruments_with_unattributed_fills,
3597 positions_with_fills,
3598 )
3599 } else {
3600 self.reconcile_position_report_netting(report, account_id)
3601 }
3602 }
3603
3604 fn reconcile_position_report_hedging(
3605 &self,
3606 report: &PositionStatusReport,
3607 account_id: AccountId,
3608 instruments_with_unattributed_fills: &IndexSet<InstrumentId>,
3609 positions_with_fills: &IndexSet<PositionId>,
3610 ) -> Option<Vec<OrderEventAny>> {
3611 let venue_position_id = report.venue_position_id?;
3612
3613 if positions_with_fills.contains(&venue_position_id) {
3615 log::debug!(
3616 "Skipping hedge position {venue_position_id} reconciliation: fills already in batch"
3617 );
3618 return None;
3619 }
3620
3621 if instruments_with_unattributed_fills.contains(&report.instrument_id) {
3624 log::debug!(
3625 "Skipping hedge position {venue_position_id} reconciliation: unattributed fills in batch"
3626 );
3627 return None;
3628 }
3629
3630 log::debug!(
3631 "Reconciling HEDGE position for {}, venue_position_id={}",
3632 report.instrument_id,
3633 venue_position_id
3634 );
3635
3636 let position = {
3637 let cache = self.cache.borrow();
3638 cache.position_owned(&venue_position_id)
3639 };
3640
3641 match position {
3642 Some(position) => {
3643 let cached_signed_qty = position.signed_decimal_qty();
3644 let venue_signed_qty = report.signed_decimal_qty;
3645
3646 if cached_signed_qty == venue_signed_qty {
3647 log::debug!(
3648 "Hedge position {venue_position_id} matches venue: qty={cached_signed_qty}"
3649 );
3650 return None;
3651 }
3652
3653 if venue_signed_qty == Decimal::ZERO && cached_signed_qty == Decimal::ZERO {
3654 return None;
3655 }
3656
3657 if !self.config.generate_missing_orders {
3658 log::error!(
3659 "Cannot reconcile {} {}: position net qty {} != reported net qty {} \
3660 and `generate_missing_orders` is disabled",
3661 report.instrument_id,
3662 venue_position_id,
3663 cached_signed_qty,
3664 venue_signed_qty
3665 );
3666 return None;
3667 }
3668
3669 self.reconcile_hedge_position_discrepancy(
3670 report,
3671 account_id,
3672 &position,
3673 cached_signed_qty,
3674 )
3675 }
3676 None => {
3677 if report.signed_decimal_qty == Decimal::ZERO {
3678 return None;
3679 }
3680
3681 if !self.config.generate_missing_orders {
3682 log::error!(
3683 "Cannot reconcile position: {venue_position_id} not found and `generate_missing_orders` is disabled"
3684 );
3685 return None;
3686 }
3687
3688 self.reconcile_missing_hedge_position(report, account_id)
3689 }
3690 }
3691 }
3692
3693 fn reconcile_hedge_position_discrepancy(
3694 &self,
3695 report: &PositionStatusReport,
3696 account_id: AccountId,
3697 position: &Position,
3698 cached_signed_qty: Decimal,
3699 ) -> Option<Vec<OrderEventAny>> {
3700 let instrument = self.get_instrument(&report.instrument_id)?;
3701 let venue_signed_qty = report.signed_decimal_qty;
3702
3703 let diff = (cached_signed_qty - venue_signed_qty).abs();
3704 let diff_qty = Quantity::from_decimal_dp(diff, instrument.size_precision()).ok()?;
3705
3706 if diff_qty.is_zero() {
3707 log::debug!(
3708 "Difference quantity rounds to zero for {}, skipping",
3709 instrument.id()
3710 );
3711 return None;
3712 }
3713
3714 let venue_position_id = report.venue_position_id?;
3715 log::warn!(
3716 "Hedge position discrepancy for {} {}: cached={}, venue={}, generating reconciliation order",
3717 report.instrument_id,
3718 venue_position_id,
3719 cached_signed_qty,
3720 venue_signed_qty
3721 );
3722
3723 let current_avg_px = if position.avg_px_open > 0.0 {
3724 Decimal::from_str(&position.avg_px_open.to_string()).ok()
3725 } else {
3726 None
3727 };
3728
3729 self.create_position_reconciliation_order(
3730 report,
3731 account_id,
3732 &instrument,
3733 cached_signed_qty,
3734 diff_qty,
3735 current_avg_px,
3736 )
3737 }
3738
3739 fn reconcile_missing_hedge_position(
3740 &self,
3741 report: &PositionStatusReport,
3742 account_id: AccountId,
3743 ) -> Option<Vec<OrderEventAny>> {
3744 let instrument = self.get_instrument(&report.instrument_id)?;
3745 let venue_signed_qty = report.signed_decimal_qty;
3746
3747 let qty = venue_signed_qty.abs();
3748 let diff_qty = Quantity::from_decimal_dp(qty, instrument.size_precision()).ok()?;
3749
3750 if diff_qty.is_zero() {
3751 return None;
3752 }
3753
3754 let venue_position_id = report.venue_position_id?;
3755 log::warn!(
3756 "Missing hedge position for {} {}: venue reports {}, generating reconciliation order",
3757 report.instrument_id,
3758 venue_position_id,
3759 venue_signed_qty
3760 );
3761
3762 self.create_position_reconciliation_order(
3763 report,
3764 account_id,
3765 &instrument,
3766 Decimal::ZERO,
3767 diff_qty,
3768 None,
3769 )
3770 }
3771
3772 fn reconcile_position_report_netting(
3773 &self,
3774 report: &PositionStatusReport,
3775 account_id: AccountId,
3776 ) -> Option<Vec<OrderEventAny>> {
3777 let instrument_id = report.instrument_id;
3778
3779 log::debug!("Reconciling NET position for {instrument_id}");
3780
3781 let instrument = self.get_instrument(&instrument_id)?;
3782
3783 let (cached_signed_qty, cached_avg_px) = {
3784 let cache = self.cache.borrow();
3785 let positions =
3786 cache.positions_open(None, Some(&instrument_id), None, Some(&account_id), None);
3787
3788 if positions.is_empty() {
3789 (Decimal::ZERO, None)
3790 } else {
3791 let mut total_signed_qty = Decimal::ZERO;
3792 let mut total_value = Decimal::ZERO;
3793 let mut total_qty = Decimal::ZERO;
3794
3795 for pos in positions {
3796 total_signed_qty += pos.signed_decimal_qty();
3797 let qty = pos.signed_decimal_qty().abs();
3798 if pos.avg_px_open > 0.0
3799 && qty > Decimal::ZERO
3800 && let Ok(avg_px) = Decimal::from_str(&pos.avg_px_open.to_string())
3801 {
3802 total_value += avg_px * qty;
3803 total_qty += qty;
3804 }
3805 }
3806
3807 let avg_px = if total_qty > Decimal::ZERO {
3808 Some(total_value / total_qty)
3809 } else {
3810 None
3811 };
3812
3813 (total_signed_qty, avg_px)
3814 }
3815 };
3816
3817 let venue_signed_qty = report.signed_decimal_qty;
3818
3819 log::debug!("venue_signed_qty={venue_signed_qty}, cached_signed_qty={cached_signed_qty}");
3820
3821 let tolerance = self.position_reconciliation_tolerance(account_id);
3822 if (cached_signed_qty - venue_signed_qty).abs() <= tolerance {
3823 log::debug!("Position quantities match for {instrument_id}, no reconciliation needed");
3824 return None;
3825 }
3826
3827 if !self.config.generate_missing_orders {
3828 log::debug!(
3829 "Discrepancy for {instrument_id} position when `generate_missing_orders` disabled, skipping"
3830 );
3831 return None;
3832 }
3833
3834 let diff = (cached_signed_qty - venue_signed_qty).abs();
3835 let diff_qty = Quantity::from_decimal_dp(diff, instrument.size_precision()).ok()?;
3836
3837 if diff_qty.is_zero() {
3838 log::debug!(
3839 "Difference quantity rounds to zero for {instrument_id}, skipping order generation"
3840 );
3841 return None;
3842 }
3843
3844 let crosses_zero = cached_signed_qty != Decimal::ZERO
3845 && venue_signed_qty != Decimal::ZERO
3846 && ((cached_signed_qty > Decimal::ZERO && venue_signed_qty < Decimal::ZERO)
3847 || (cached_signed_qty < Decimal::ZERO && venue_signed_qty > Decimal::ZERO));
3848
3849 if crosses_zero {
3850 let ts_now = self.clock.borrow().timestamp_ns();
3851 return self.reconcile_cross_zero_position(
3852 &instrument,
3853 account_id,
3854 instrument_id,
3855 cached_signed_qty,
3856 cached_avg_px,
3857 venue_signed_qty,
3858 report.avg_px_open,
3859 ts_now,
3860 report.ts_last,
3861 );
3862 }
3863
3864 if cached_signed_qty == Decimal::ZERO {
3865 return self.create_position_from_report(report, account_id, &instrument);
3866 }
3867
3868 self.create_position_reconciliation_order(
3869 report,
3870 account_id,
3871 &instrument,
3872 cached_signed_qty,
3873 diff_qty,
3874 cached_avg_px,
3875 )
3876 }
3877
3878 fn create_position_reconciliation_order(
3879 &self,
3880 report: &PositionStatusReport,
3881 account_id: AccountId,
3882 instrument: &InstrumentAny,
3883 cached_signed_qty: Decimal,
3884 diff_qty: Quantity,
3885 current_avg_px: Option<Decimal>,
3886 ) -> Option<Vec<OrderEventAny>> {
3887 let venue_signed_qty = report.signed_decimal_qty;
3888 let instrument_id = report.instrument_id;
3889
3890 let order_side = if venue_signed_qty > cached_signed_qty {
3891 OrderSide::Buy
3892 } else {
3893 OrderSide::Sell
3894 };
3895
3896 let reconciliation_px = calculate_reconciliation_price(
3897 cached_signed_qty,
3898 current_avg_px,
3899 venue_signed_qty,
3900 report.avg_px_open,
3901 );
3902
3903 let fill_px = reconciliation_px
3904 .or(report.avg_px_open)
3905 .or(current_avg_px)?;
3906
3907 let ts_now = self.clock.borrow().timestamp_ns();
3908 let fill_price = Price::from_decimal_dp(fill_px, instrument.price_precision()).ok();
3909 let venue_order_id = create_position_reconciliation_venue_order_id(
3910 account_id,
3911 instrument_id,
3912 order_side,
3913 OrderType::Market,
3914 diff_qty,
3915 fill_price,
3916 report.venue_position_id,
3917 None,
3918 report.ts_last,
3919 );
3920
3921 let mut order_report = OrderStatusReport::new(
3922 account_id,
3923 instrument_id,
3924 None,
3925 venue_order_id,
3926 order_side,
3927 OrderType::Market,
3928 TimeInForce::Gtc,
3929 OrderStatus::Filled,
3930 diff_qty,
3931 diff_qty,
3932 ts_now,
3933 ts_now,
3934 ts_now,
3935 None,
3936 )
3937 .with_avg_px(fill_px);
3938
3939 if let Some(venue_position_id) = report.venue_position_id {
3940 order_report = order_report.with_venue_position_id(venue_position_id);
3941 }
3942
3943 log::info!(
3944 color = LogColor::Blue as u8;
3945 "Generating reconciliation order for {instrument_id}: side={order_side:?}, qty={diff_qty}, px={fill_px}",
3946 );
3947
3948 let (events, _) = self.handle_external_order(
3949 &order_report,
3950 account_id,
3951 instrument,
3952 &[],
3953 true,
3954 None,
3955 None,
3956 );
3957 Some(events)
3958 }
3959
3960 fn reconcile_order_report(
3961 &self,
3962 order: &OrderAny,
3963 report: &OrderStatusReport,
3964 instrument: Option<&InstrumentAny>,
3965 commission_client: Option<&dyn ExecutionClient>,
3966 ) -> anyhow::Result<Option<OrderEventAny>> {
3967 let ts_now = self.clock.borrow().timestamp_ns();
3968 let commission = if matches!(
3969 report.order_status,
3970 OrderStatus::PartiallyFilled | OrderStatus::Filled
3971 ) && report.filled_qty > order.filled_qty()
3972 {
3973 let Some(instrument) = instrument else {
3974 return Ok(reconcile_order_report_with_commission(
3975 order, report, None, ts_now, None,
3976 ));
3977 };
3978 let fill_qty = report.filled_qty - order.filled_qty();
3979 Self::resolve_inferred_fill_commission(
3980 commission_client,
3981 instrument,
3982 fill_qty,
3983 incremental_inferred_fill_price_and_liquidity(order, report, instrument),
3984 )?
3985 } else {
3986 None
3987 };
3988
3989 Ok(reconcile_order_report_with_commission(
3990 order, report, instrument, ts_now, commission,
3991 ))
3992 }
3993
3994 fn reconcile_order_with_fills(
3996 &mut self,
3997 order: &OrderAny,
3998 report: &OrderStatusReport,
3999 fills: &[&FillReport],
4000 instrument: Option<&InstrumentAny>,
4001 fill_queue: &mut ReconciliationFillQueue,
4002 commission_client: Option<&dyn ExecutionClient>,
4003 ) -> Vec<OrderEventAny> {
4004 let mut events = Vec::new();
4005 let mut working = order.clone();
4006 let mut sorted_fills: Vec<&FillReport> = fills.to_vec();
4007 sorted_fills.sort_by_key(|f| f.ts_event);
4008
4009 let ts_now = self.clock.borrow().timestamp_ns();
4010
4011 if matches!(
4012 report.order_status,
4013 OrderStatus::Canceled | OrderStatus::Expired
4014 ) && report.ts_triggered.is_some()
4015 && working.status() != OrderStatus::Triggered
4016 && TRIGGERABLE_ORDER_TYPES.contains(&working.order_type())
4017 {
4018 let triggered = create_reconciliation_triggered(&working, report, ts_now);
4019 if working.apply(triggered.clone()).is_ok() {
4020 events.push(triggered);
4021 }
4022 }
4023
4024 let requires_snapshot_projection = !sorted_fills.is_empty()
4025 || report.order_status == OrderStatus::Voided
4026 || report.filled_qty < working.filled_qty();
4027 if !requires_snapshot_projection {
4028 match self.reconcile_order_report(&working, report, instrument, commission_client) {
4029 Ok(Some(event)) => events.push(event),
4030 Ok(None) => {}
4031 Err(e) => log::error!(
4032 "Deferring inferred fill for {}: venue commission calculation failed: {e}",
4033 order.client_order_id(),
4034 ),
4035 }
4036 return events;
4037 }
4038
4039 for event in generate_reconciliation_order_pre_fill_events(&working, report, ts_now) {
4040 if let Err(e) = working.apply(event.clone()) {
4041 log::warn!(
4042 "Cannot project reconciliation event for {}: {e}",
4043 order.client_order_id()
4044 );
4045 return events;
4046 }
4047 events.push(event);
4048 }
4049
4050 if let Some(inst) = instrument {
4051 for fill in sorted_fills {
4052 let Some((event, fill_key)) =
4053 self.create_order_fill(&working, fill, inst, &fill_queue.pending_fill_keys)
4054 else {
4055 continue;
4056 };
4057
4058 if let Err(e) = working.apply(event.clone()) {
4059 if let OrderEventAny::Filled(fill) = &event
4060 && self.is_fill_applied(fill, fill_key)
4061 {
4062 self.processed_fills.mark(fill_key);
4063 } else {
4064 log::warn!(
4065 "Cannot project reconciliation fill for {}: {e}",
4066 order.client_order_id()
4067 );
4068 }
4069 return events;
4070 }
4071 fill_queue.push(&mut events, event, fill_key);
4072 }
4073 }
4074
4075 let commission = if report.filled_qty > working.filled_qty()
4076 && let Some(instrument) = instrument
4077 {
4078 let fill_qty = report.filled_qty - working.filled_qty();
4079
4080 match Self::resolve_inferred_fill_commission(
4081 commission_client,
4082 instrument,
4083 fill_qty,
4084 incremental_inferred_fill_price_and_liquidity(&working, report, instrument),
4085 ) {
4086 Ok(commission) => commission,
4087 Err(e) => {
4088 log::error!(
4089 "Deferring inferred fill for {}: venue commission calculation failed: {e}",
4090 order.client_order_id(),
4091 );
4092 return events;
4093 }
4094 }
4095 } else {
4096 None
4097 };
4098
4099 for event in generate_reconciliation_order_snapshot_events_with_commission(
4100 &working, report, instrument, ts_now, commission,
4101 ) {
4102 if let Err(e) = working.apply(event.clone()) {
4103 log::warn!(
4104 "Cannot project reconciliation snapshot event for {}: {e}",
4105 order.client_order_id()
4106 );
4107 break;
4108 }
4109 events.push(event);
4110 }
4111
4112 events
4113 }
4114
4115 fn resolve_inferred_fill_commission(
4116 client: Option<&dyn ExecutionClient>,
4117 instrument: &InstrumentAny,
4118 fill_qty: Quantity,
4119 price_and_liquidity: Option<(Price, LiquiditySide)>,
4120 ) -> anyhow::Result<Option<Money>> {
4121 let Some(client) = client else {
4122 anyhow::bail!("responsible execution client is unavailable");
4123 };
4124 let Some((last_px, liquidity_side)) = price_and_liquidity else {
4125 return Ok(None);
4126 };
4127
4128 client.calculate_commission(instrument, fill_qty, last_px, liquidity_side)
4129 }
4130
4131 #[expect(clippy::too_many_arguments)]
4132 fn handle_external_order(
4133 &self,
4134 report: &OrderStatusReport,
4135 account_id: AccountId,
4136 instrument: &InstrumentAny,
4137 fills: &[&FillReport],
4138 is_synthetic: bool,
4139 fill_queue: Option<&mut ReconciliationFillQueue>,
4140 commission_client: Option<&dyn ExecutionClient>,
4141 ) -> (Vec<OrderEventAny>, Option<ExternalOrderMetadata>) {
4142 let (strategy_id, tags) =
4143 if let Some(claimed_strategy) = self.external_order_claims.get(&report.instrument_id) {
4144 let order_id = report
4145 .client_order_id
4146 .map_or_else(|| report.venue_order_id.to_string(), |id| id.to_string());
4147 log::info!(
4148 color = LogColor::Blue as u8;
4149 "External order {} for {} claimed by strategy {}",
4150 order_id,
4151 report.instrument_id,
4152 claimed_strategy,
4153 );
4154 (*claimed_strategy, None)
4155 } else {
4156 let tag = if is_synthetic {
4158 *TAG_RECONCILIATION
4159 } else {
4160 *TAG_VENUE
4161 };
4162 (StrategyId::from("EXTERNAL"), Some(vec![tag]))
4163 };
4164
4165 if self.config.filter_unclaimed_external && !is_synthetic {
4167 return (Vec::new(), None);
4168 }
4169
4170 let client_order_id = report
4171 .client_order_id
4172 .unwrap_or_else(|| ClientOrderId::from(report.venue_order_id.as_str()));
4173
4174 if !report.quantity.is_positive() {
4175 log::error!(
4176 "Skipping external order {} ({}) for {}: non-positive quantity in report {:?}",
4177 client_order_id,
4178 report.venue_order_id,
4179 report.instrument_id,
4180 report,
4181 );
4182 return (Vec::new(), None);
4183 }
4184
4185 let ts_now = self.clock.borrow().timestamp_ns();
4186
4187 let initialized = match OrderInitialized::new_checked(
4188 self.config.trader_id,
4189 strategy_id,
4190 report.instrument_id,
4191 client_order_id,
4192 report.order_side,
4193 report.order_type,
4194 report.quantity,
4195 report.time_in_force,
4196 report.post_only,
4197 report.reduce_only,
4198 false, true, UUID4::new(),
4201 ts_now,
4202 ts_now,
4203 report.price,
4204 report.activation_price,
4205 report.trigger_price,
4206 report.trigger_type,
4207 report.limit_offset,
4208 report.trailing_offset,
4209 Some(report.trailing_offset_type),
4210 report.expire_time,
4211 report.display_qty,
4212 None, None, Some(report.contingency_type),
4215 report.order_list_id,
4216 report.linked_order_ids.clone(),
4217 report.parent_order_id,
4218 None, None, None, tags,
4222 ) {
4223 Ok(initialized) => initialized,
4224 Err(e) => {
4225 log::error!("Failed to create order from report: {e}");
4226 return (Vec::new(), None);
4227 }
4228 };
4229
4230 let initialized = OrderEventAny::Initialized(initialized);
4231 let order = match OrderAny::from_events(vec![initialized.clone()]) {
4232 Ok(order) => order,
4233 Err(e) => {
4234 log::error!("Failed to create order from report: {e}");
4235 return (Vec::new(), None);
4236 }
4237 };
4238
4239 let replace_inferred_fill = !fills.is_empty()
4240 && matches!(
4241 report.order_status,
4242 OrderStatus::Canceled
4243 | OrderStatus::Expired
4244 | OrderStatus::Filled
4245 | OrderStatus::PartiallyFilled
4246 );
4247 let mut prepared_fills = Vec::new();
4248 let mut prepared_fill_keys = fill_queue
4249 .as_deref()
4250 .map(|queue| queue.pending_fill_keys.clone())
4251 .unwrap_or_default();
4252 let mut real_fill_total = Decimal::ZERO;
4253
4254 if replace_inferred_fill {
4255 let mut sorted_fills: Vec<&FillReport> = fills.to_vec();
4256 sorted_fills.sort_by_key(|fill| fill.ts_event);
4257 if fill_queue.is_none() {
4258 log::error!(
4259 "Cannot reconcile external order {client_order_id}: fill queue is unavailable"
4260 );
4261 return (Vec::new(), None);
4262 }
4263
4264 for fill in sorted_fills {
4265 if let Some((fill_event, fill_key)) =
4266 self.create_order_fill(&order, fill, instrument, &prepared_fill_keys)
4267 {
4268 real_fill_total += fill.last_qty.as_decimal();
4269 prepared_fill_keys.insert(fill_key);
4270 prepared_fills.push((fill_event, fill_key));
4271 }
4272 }
4273 }
4274
4275 let report_filled = report.filled_qty.as_decimal();
4276 let inferred_qty = if report_filled.is_zero() {
4277 None
4278 } else if replace_inferred_fill {
4279 if real_fill_total < report_filled {
4280 match Quantity::from_decimal_dp(
4281 report_filled - real_fill_total,
4282 instrument.size_precision(),
4283 ) {
4284 Ok(quantity) => Some(quantity),
4285 Err(e) => {
4286 log::error!(
4287 "Cannot reconcile external order {client_order_id}: residual fill quantity is invalid: {e}"
4288 );
4289 return (Vec::new(), None);
4290 }
4291 }
4292 } else {
4293 None
4294 }
4295 } else if matches!(
4296 report.order_status,
4297 OrderStatus::PartiallyFilled
4298 | OrderStatus::Filled
4299 | OrderStatus::Canceled
4300 | OrderStatus::Expired
4301 | OrderStatus::Voided
4302 ) {
4303 Some(report.filled_qty)
4304 } else {
4305 None
4306 };
4307
4308 let inferred_commission = if is_synthetic {
4309 None
4310 } else if let Some(inferred_qty) = inferred_qty {
4311 match Self::resolve_inferred_fill_commission(
4312 commission_client,
4313 instrument,
4314 inferred_qty,
4315 inferred_fill_price_and_liquidity(&order, report, instrument),
4316 ) {
4317 Ok(commission) => commission,
4318 Err(e) => {
4319 log::error!(
4320 "Deferring external order {client_order_id}: venue commission calculation failed: {e}"
4321 );
4322 return (Vec::new(), None);
4323 }
4324 }
4325 } else {
4326 None
4327 };
4328
4329 {
4330 let mut cache = self.cache.borrow_mut();
4331 let source_client_id = if is_synthetic {
4332 None
4333 } else {
4334 commission_client.map(ExecutionClient::client_id)
4335 };
4336
4337 if let Err(e) = cache.add_order(order.clone(), None, source_client_id, false) {
4338 match cache.order(&client_order_id) {
4342 Some(existing) if is_synthetic && existing.is_closed() => {
4343 log::debug!(
4344 "Skipping synthetic reconciliation order {client_order_id} for {}: \
4345 replay deduped (cached status={:?})",
4346 report.instrument_id,
4347 existing.status(),
4348 );
4349 }
4350 Some(existing) if is_synthetic => {
4351 log::warn!(
4352 "Synthetic reconciliation order {client_order_id} for {} exists in \
4353 cache in non-terminal state {:?}; fill not regenerated",
4354 report.instrument_id,
4355 existing.status(),
4356 );
4357 }
4358 _ => {
4359 log::error!("Failed to add external order to cache: {e}");
4360 }
4361 }
4362 return (Vec::new(), None);
4363 }
4364
4365 if let Err(e) = cache.index_venue_order_id(&client_order_id, &report.venue_order_id) {
4366 log::warn!("Failed to index venue order ID: {e}");
4367 }
4368 }
4369
4370 Self::publish_order_event(&initialized);
4371
4372 log::info!(
4373 color = LogColor::Blue as u8;
4374 "Created external order {} ({}) for {} [{}]",
4375 client_order_id,
4376 report.venue_order_id,
4377 report.instrument_id,
4378 report.order_status,
4379 );
4380
4381 let ts_now = self.clock.borrow().timestamp_ns();
4382 let mut order_events = generate_external_order_status_events_with_commission(
4383 &order,
4384 report,
4385 &account_id,
4386 instrument,
4387 ts_now,
4388 inferred_commission,
4389 );
4390
4391 if replace_inferred_fill {
4392 let terminal_event = if order_events.last().is_some_and(|event| {
4393 matches!(
4394 event,
4395 OrderEventAny::Canceled(_) | OrderEventAny::Expired(_),
4396 )
4397 }) {
4398 order_events.pop()
4399 } else {
4400 None
4401 };
4402
4403 if order_events
4404 .last()
4405 .is_some_and(|event| matches!(event, OrderEventAny::Filled(_)))
4406 {
4407 order_events.pop();
4408 }
4409
4410 let fill_queue =
4411 fill_queue.expect("fill queue availability was checked before cache mutation");
4412 for (fill_event, fill_key) in prepared_fills {
4413 fill_queue.push(&mut order_events, fill_event, fill_key);
4414 }
4415
4416 if let Some(inferred_qty) = inferred_qty
4417 && let Some(inferred_fill) = create_inferred_fill_for_qty(
4418 &order,
4419 report,
4420 &account_id,
4421 instrument,
4422 inferred_qty,
4423 ts_now,
4424 inferred_commission,
4425 )
4426 {
4427 order_events.push(inferred_fill);
4428 }
4429
4430 if let Some(event) = terminal_event {
4431 order_events.push(event);
4432 }
4433 }
4434
4435 let metadata = ExternalOrderMetadata {
4436 client_order_id,
4437 venue_order_id: report.venue_order_id,
4438 instrument_id: report.instrument_id,
4439 strategy_id,
4440 ts_init: ts_now,
4441 };
4442
4443 (order_events, Some(metadata))
4444 }
4445
4446 fn publish_order_event(event: &OrderEventAny) {
4447 let topic = switchboard::get_event_order_topic(event.strategy_id());
4448 msgbus::publish_order_event(topic, event);
4449 }
4450
4451 fn adjust_mass_status_fills(
4456 &self,
4457 mass_status: &ExecutionMassStatus,
4458 ) -> (
4459 IndexMap<VenueOrderId, OrderStatusReport>,
4460 IndexMap<VenueOrderId, Vec<FillReport>>,
4461 ) {
4462 let mut final_orders: IndexMap<VenueOrderId, OrderStatusReport> =
4463 mass_status.order_reports();
4464 let mut final_fills: IndexMap<VenueOrderId, Vec<FillReport>> = mass_status.fill_reports();
4465
4466 if mass_status.lookback_start().is_some() {
4467 return (final_orders, final_fills);
4468 }
4469
4470 let mut instruments_to_adjust = Vec::new();
4471
4472 for (instrument_id, position_reports) in mass_status.position_reports() {
4473 if !self.should_reconcile_instrument(&instrument_id) {
4474 log::debug!(
4475 "Skipping fill adjustment for {instrument_id}: not in reconciliation_instrument_ids"
4476 );
4477 continue;
4478 }
4479
4480 let is_hedge_mode = position_reports
4483 .iter()
4484 .any(|r| r.venue_position_id.is_some());
4485
4486 if is_hedge_mode {
4487 log::debug!(
4488 "Skipping fill adjustment for {instrument_id}: hedge mode (has venue_position_id)"
4489 );
4490 continue;
4491 }
4492
4493 let has_retained_position = {
4494 let cache = self.cache.borrow();
4495 !cache
4496 .positions_open(
4497 None,
4498 Some(&instrument_id),
4499 None,
4500 Some(&mass_status.account_id),
4501 None,
4502 )
4503 .is_empty()
4504 };
4505
4506 if has_retained_position {
4507 log::debug!(
4508 "Skipping fill adjustment for {instrument_id}: retained open position in cache"
4509 );
4510 continue;
4511 }
4512
4513 if let Some(instrument) = self.get_instrument(&instrument_id) {
4514 instruments_to_adjust.push(instrument);
4515 } else {
4516 log::debug!(
4517 "Skipping fill adjustment for {instrument_id}: instrument not found in cache"
4518 );
4519 }
4520 }
4521
4522 if instruments_to_adjust.is_empty() {
4523 return (final_orders, final_fills);
4524 }
4525
4526 log_info!(
4527 "Adjusting fills for {} instrument(s) with position reports",
4528 instruments_to_adjust.len(),
4529 color = LogColor::Blue
4530 );
4531
4532 for instrument in &instruments_to_adjust {
4533 let instrument_id = instrument.id();
4534
4535 match process_mass_status_for_reconciliation(mass_status, instrument, None) {
4536 Ok(result) => {
4537 final_orders.retain(|_, order| order.instrument_id != instrument_id);
4538 final_fills.retain(|_, fills| {
4539 fills
4540 .first()
4541 .is_none_or(|f| f.instrument_id != instrument_id)
4542 });
4543
4544 for (venue_order_id, order) in result.orders {
4545 final_orders.insert(venue_order_id, order);
4546 }
4547
4548 for (venue_order_id, fills) in result.fills {
4549 final_fills.insert(venue_order_id, fills);
4550 }
4551 }
4552 Err(e) => {
4553 log::warn!("Failed to adjust fills for {instrument_id}: {e}");
4554 }
4555 }
4556 }
4557
4558 log_info!(
4559 "After adjustment: {} order(s), {} fill group(s)",
4560 final_orders.len(),
4561 final_fills.len(),
4562 color = LogColor::Blue
4563 );
4564
4565 (final_orders, final_fills)
4566 }
4567
4568 fn deduplicate_order_reports<'a>(
4573 reports: impl Iterator<Item = &'a OrderStatusReport>,
4574 ) -> IndexMap<VenueOrderId, &'a OrderStatusReport> {
4575 let mut best_reports: IndexMap<VenueOrderId, &'a OrderStatusReport> = IndexMap::new();
4576
4577 for report in reports {
4578 let dominated = best_reports
4579 .get(&report.venue_order_id)
4580 .is_some_and(|existing| Self::is_more_advanced(existing, report));
4581
4582 if !dominated {
4583 best_reports.insert(report.venue_order_id, report);
4584 }
4585 }
4586
4587 best_reports
4588 }
4589
4590 fn is_more_advanced(a: &OrderStatusReport, b: &OrderStatusReport) -> bool {
4591 if a.filled_qty > b.filled_qty {
4592 return true;
4593 }
4594
4595 if a.filled_qty < b.filled_qty {
4596 return false;
4597 }
4598
4599 Self::status_priority(a.order_status) > Self::status_priority(b.order_status)
4601 }
4602
4603 const fn status_priority(status: OrderStatus) -> u8 {
4604 match status {
4605 OrderStatus::Initialized | OrderStatus::Submitted | OrderStatus::Emulated => 0,
4606 OrderStatus::Released | OrderStatus::Denied => 1,
4607 OrderStatus::Accepted | OrderStatus::PendingUpdate | OrderStatus::PendingCancel => 2,
4608 OrderStatus::Triggered => 3,
4609 OrderStatus::PartiallyFilled => 4,
4610 OrderStatus::Canceled | OrderStatus::Expired | OrderStatus::Rejected => 5,
4611 OrderStatus::Filled | OrderStatus::Voided => 6,
4612 }
4613 }
4614
4615 fn is_exact_order_match(order: &OrderAny, report: &OrderStatusReport) -> bool {
4616 order.status() == report.order_status
4617 && order.filled_qty() == report.filled_qty
4618 && !should_reconciliation_update(order, report)
4619 }
4620
4621 fn is_fill_applied(&self, fill: &OrderFilled, fill_key: FillKey) -> bool {
4622 self.get_order(fill.client_order_id)
4623 .or_else(|| self.get_order_by_venue_order_id(fill.venue_order_id))
4624 .is_some_and(|order| {
4625 order.account_id() == Some(fill_key.0)
4626 && order.instrument_id() == fill_key.1
4627 && order.trade_ids().contains(&&fill_key.2)
4628 })
4629 }
4630
4631 fn create_order_fill(
4632 &self,
4633 order: &OrderAny,
4634 fill: &FillReport,
4635 instrument: &InstrumentAny,
4636 pending_fill_keys: &IndexSet<FillKey>,
4637 ) -> Option<(OrderEventAny, FillKey)> {
4638 let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
4639 if self.processed_fills.contains_key(&fill_key) || pending_fill_keys.contains(&fill_key) {
4640 return None;
4641 }
4642
4643 let event = OrderEventAny::Filled(OrderFilled::new(
4644 order.trader_id(),
4645 order.strategy_id(),
4646 order.instrument_id(),
4647 order.client_order_id(),
4648 fill.venue_order_id,
4649 fill.account_id,
4650 fill.trade_id,
4651 fill.order_side,
4652 order.order_type(),
4653 fill.last_qty,
4654 fill.last_px,
4655 instrument.quote_currency(),
4656 fill.liquidity_side,
4657 fill.report_id,
4658 fill.ts_event,
4659 self.clock.borrow().timestamp_ns(),
4660 false,
4661 fill.venue_position_id,
4662 Some(fill.commission),
4663 None,
4664 ));
4665
4666 Some((event, fill_key))
4667 }
4668}
4669
4670pub(crate) async fn request_targeted_order_reports(
4671 clients: &[&dyn ExecutionClient],
4672 queries: Vec<TargetedOrderQuery>,
4673 query_delay: Duration,
4674) -> Vec<TargetedOrderReportResult> {
4675 let mut results = Vec::with_capacity(queries.len());
4676 let mut request_count = 0usize;
4677
4678 for query in queries {
4679 let mut report = None;
4680 let mut report_client_id = None;
4681 let mut coverage_complete = true;
4682
4683 for client_id in &query.responsible_clients {
4684 let client_id = *client_id;
4685 let Some(client) = clients
4686 .iter()
4687 .find(|client| client.client_id() == client_id)
4688 else {
4689 coverage_complete = false;
4690 log::warn!(
4691 "Cannot run targeted order status query for {}: execution client {client_id} is unavailable",
4692 query.client_order_id,
4693 );
4694 continue;
4695 };
4696
4697 if request_count > 0 && !query_delay.is_zero() {
4698 dst::time::sleep(query_delay).await;
4699 }
4700 request_count += 1;
4701
4702 match client.generate_order_status_report(&query.command).await {
4703 Ok(Some(candidate)) if targeted_report_matches(&query, &candidate) => {
4704 report = Some(candidate);
4705 report_client_id = Some(client_id);
4706 break;
4707 }
4708 Ok(Some(candidate)) => {
4709 coverage_complete = false;
4710 log::warn!(
4711 "Ignoring mismatched targeted order status report from {client_id} for {}: client_order_id={:?}, venue_order_id={}, instrument_id={}",
4712 query.client_order_id,
4713 candidate.client_order_id,
4714 candidate.venue_order_id,
4715 candidate.instrument_id,
4716 );
4717 }
4718 Ok(None) => {}
4719 Err(e) => {
4720 coverage_complete = false;
4721 log::warn!(
4722 "Failed targeted order status query from {client_id} for {}: {e}",
4723 query.client_order_id,
4724 );
4725 }
4726 }
4727 }
4728
4729 results.push(TargetedOrderReportResult {
4730 client_order_id: query.client_order_id,
4731 client_id: report_client_id,
4732 report,
4733 coverage_complete,
4734 });
4735 }
4736
4737 results
4738}
4739
4740fn targeted_report_matches(query: &TargetedOrderQuery, report: &OrderStatusReport) -> bool {
4741 let instrument_matches = query
4742 .command
4743 .instrument_id
4744 .is_none_or(|instrument_id| report.instrument_id == instrument_id);
4745 let order_matches = report.client_order_id == Some(query.client_order_id)
4746 || query
4747 .command
4748 .venue_order_id
4749 .is_some_and(|venue_order_id| report.venue_order_id == venue_order_id);
4750
4751 instrument_matches && order_matches
4752}
4753
4754#[cfg(test)]
4755mod tests {
4756 use nautilus_common::clock::TestClock;
4757 use nautilus_core::{Params, datetime::NANOSECONDS_IN_SECOND};
4758 use nautilus_execution::reconciliation::generate_reconciliation_order_events;
4759 use nautilus_model::{
4760 accounts::AccountAny,
4761 enums::{LiquiditySide, OmsType, PositionSideSpecified},
4762 events::order::spec::{OrderPendingUpdateSpec, OrderUpdatedSpec},
4763 identifiers::Venue,
4764 instruments::{
4765 Instrument,
4766 stubs::{crypto_perpetual_ethusdt, xbtusd_bitmex},
4767 },
4768 orders::{OrderTestBuilder, stubs::TestOrderEventStubs},
4769 types::{AccountBalance, Currency, MarginBalance, Money},
4770 };
4771 use rstest::rstest;
4772 use rust_decimal_macros::dec;
4773
4774 use super::*;
4775
4776 #[derive(Clone)]
4777 enum CommissionOutcome {
4778 Value(Money),
4779 NoOverride,
4780 Failure,
4781 }
4782
4783 struct CommissionStubClient {
4784 outcome: CommissionOutcome,
4785 seen: RefCell<Option<(Quantity, Price, LiquiditySide)>>,
4786 }
4787
4788 impl CommissionStubClient {
4789 fn new(outcome: CommissionOutcome) -> Self {
4790 Self {
4791 outcome,
4792 seen: RefCell::new(None),
4793 }
4794 }
4795 }
4796
4797 #[async_trait::async_trait(?Send)]
4798 impl ExecutionClient for CommissionStubClient {
4799 fn is_connected(&self) -> bool {
4800 true
4801 }
4802
4803 fn client_id(&self) -> ClientId {
4804 ClientId::from("STUB")
4805 }
4806
4807 fn account_id(&self) -> AccountId {
4808 AccountId::from("STUB-001")
4809 }
4810
4811 fn venue(&self) -> Venue {
4812 Venue::from("STUB")
4813 }
4814
4815 fn oms_type(&self) -> OmsType {
4816 OmsType::Netting
4817 }
4818
4819 fn get_account(&self) -> Option<AccountAny> {
4820 None
4821 }
4822
4823 fn generate_account_state(
4824 &self,
4825 _balances: Vec<AccountBalance>,
4826 _margins: Vec<MarginBalance>,
4827 _reported: bool,
4828 _ts_event: UnixNanos,
4829 _info: Option<Params>,
4830 ) -> anyhow::Result<()> {
4831 Ok(())
4832 }
4833
4834 fn start(&mut self) -> anyhow::Result<()> {
4835 Ok(())
4836 }
4837
4838 fn stop(&mut self) -> anyhow::Result<()> {
4839 Ok(())
4840 }
4841
4842 fn calculate_commission(
4843 &self,
4844 _instrument: &InstrumentAny,
4845 last_qty: Quantity,
4846 last_px: Price,
4847 liquidity_side: LiquiditySide,
4848 ) -> anyhow::Result<Option<Money>> {
4849 *self.seen.borrow_mut() = Some((last_qty, last_px, liquidity_side));
4850
4851 match &self.outcome {
4852 CommissionOutcome::Value(money) => Ok(Some(*money)),
4853 CommissionOutcome::NoOverride => Ok(None),
4854 CommissionOutcome::Failure => {
4855 anyhow::bail!("commission is not representable as Money")
4856 }
4857 }
4858 }
4859 }
4860
4861 fn commission_fixtures() -> (OrderAny, OrderStatusReport, InstrumentAny) {
4862 let instrument = crypto_perpetual_ethusdt();
4863 let order = OrderTestBuilder::new(OrderType::Limit)
4864 .instrument_id(instrument.id())
4865 .side(OrderSide::Buy)
4866 .quantity(Quantity::from("10.0"))
4867 .price(Price::from("100.00"))
4868 .build();
4869 let report = OrderStatusReport::new(
4870 AccountId::from("STUB-001"),
4871 instrument.id(),
4872 Some(order.client_order_id()),
4873 VenueOrderId::from("V-1"),
4874 OrderSide::Buy,
4875 OrderType::Limit,
4876 TimeInForce::Gtc,
4877 OrderStatus::Filled,
4878 Quantity::from("10.0"),
4879 Quantity::from("10.0"),
4880 UnixNanos::from(1),
4881 UnixNanos::from(1),
4882 UnixNanos::from(1),
4883 None,
4884 )
4885 .with_avg_px(dec!(100.0));
4886
4887 (order, report, InstrumentAny::CryptoPerpetual(instrument))
4888 }
4889
4890 fn cached_commission_fixtures() -> (
4891 ExecutionManager,
4892 Rc<RefCell<Cache>>,
4893 OrderAny,
4894 OrderStatusReport,
4895 InstrumentAny,
4896 ) {
4897 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
4898 let clock = Rc::new(RefCell::new(TestClock::new()));
4899 let cache = Rc::new(RefCell::new(Cache::default()));
4900 cache
4901 .borrow_mut()
4902 .add_instrument(instrument.clone())
4903 .expect("instrument is cacheable");
4904 let client_order_id = ClientOrderId::from("O-COMMISSION-CACHED");
4905 let venue_order_id = VenueOrderId::from("V-COMMISSION-CACHED");
4906 insert_accepted_limit_order(
4907 &cache,
4908 client_order_id,
4909 venue_order_id,
4910 instrument.id(),
4911 ClientId::from("STUB"),
4912 );
4913 let order = cache
4914 .borrow()
4915 .order_owned(&client_order_id)
4916 .expect("accepted order is cached");
4917 let report = OrderStatusReport::new(
4918 AccountId::from("TEST-001"),
4919 instrument.id(),
4920 Some(client_order_id),
4921 venue_order_id,
4922 OrderSide::Buy,
4923 OrderType::Limit,
4924 TimeInForce::Gtc,
4925 OrderStatus::Filled,
4926 Quantity::from("10.0"),
4927 Quantity::from("10.0"),
4928 UnixNanos::from(1),
4929 UnixNanos::from(1),
4930 UnixNanos::from(1),
4931 None,
4932 )
4933 .with_avg_px(dec!(100.0));
4934 let manager =
4935 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default());
4936
4937 (manager, cache, order, report, instrument)
4938 }
4939
4940 #[rstest]
4941 fn test_resolve_inferred_fill_commission_without_client_fails_closed() {
4942 let (order, report, instrument) = commission_fixtures();
4943
4944 let error = ExecutionManager::resolve_inferred_fill_commission(
4945 None,
4946 &instrument,
4947 Quantity::from("5.0"),
4948 inferred_fill_price_and_liquidity(&order, &report, &instrument),
4949 )
4950 .expect_err("a missing responsible client must defer the fill");
4951
4952 assert_eq!(
4953 error.to_string(),
4954 "responsible execution client is unavailable"
4955 );
4956 }
4957
4958 #[rstest]
4959 fn test_resolve_inferred_fill_commission_without_price_uses_generic_path() {
4960 let instrument = crypto_perpetual_ethusdt();
4961 let order = OrderTestBuilder::new(OrderType::Market)
4962 .instrument_id(instrument.id())
4963 .side(OrderSide::Buy)
4964 .quantity(Quantity::from("10.0"))
4965 .build();
4966 let report = OrderStatusReport::new(
4967 AccountId::from("STUB-001"),
4968 instrument.id(),
4969 Some(order.client_order_id()),
4970 VenueOrderId::from("V-1"),
4971 OrderSide::Buy,
4972 OrderType::Market,
4973 TimeInForce::Gtc,
4974 OrderStatus::Filled,
4975 Quantity::from("10.0"),
4976 Quantity::from("10.0"),
4977 UnixNanos::from(1),
4978 UnixNanos::from(1),
4979 UnixNanos::from(1),
4980 None,
4981 );
4982 let client =
4983 CommissionStubClient::new(CommissionOutcome::Value(Money::new(1.0, Currency::USDT())));
4984 let instrument = InstrumentAny::CryptoPerpetual(instrument);
4985
4986 let commission = ExecutionManager::resolve_inferred_fill_commission(
4987 Some(&client),
4988 &instrument,
4989 Quantity::from("5.0"),
4990 inferred_fill_price_and_liquidity(&order, &report, &instrument),
4991 )
4992 .expect("an unresolvable price is not a failure");
4993
4994 assert_eq!(commission, None, "no price means no venue commission");
4995 }
4996
4997 #[rstest]
4998 fn test_resolve_inferred_fill_commission_returns_venue_value() {
4999 let (order, report, instrument) = commission_fixtures();
5000 let expected = Money::new(2.5, Currency::USDT());
5001 let client = CommissionStubClient::new(CommissionOutcome::Value(expected));
5002
5003 let commission = ExecutionManager::resolve_inferred_fill_commission(
5004 Some(&client),
5005 &instrument,
5006 Quantity::from("5.0"),
5007 inferred_fill_price_and_liquidity(&order, &report, &instrument),
5008 )
5009 .expect("a representable commission succeeds");
5010
5011 assert_eq!(commission, Some(expected));
5012 assert_eq!(
5013 *client.seen.borrow(),
5014 Some((
5015 Quantity::from("5.0"),
5016 Price::from("100.00"),
5017 LiquiditySide::NoLiquiditySide,
5018 )),
5019 "the resolver passes the inferred fill quantity, resolved price, and liquidity side"
5020 );
5021 }
5022
5023 #[rstest]
5024 fn test_resolve_inferred_fill_commission_honors_no_override() {
5025 let (order, report, instrument) = commission_fixtures();
5026 let client = CommissionStubClient::new(CommissionOutcome::NoOverride);
5027
5028 let commission = ExecutionManager::resolve_inferred_fill_commission(
5029 Some(&client),
5030 &instrument,
5031 Quantity::from("5.0"),
5032 inferred_fill_price_and_liquidity(&order, &report, &instrument),
5033 )
5034 .expect("no override is not a failure");
5035
5036 assert_eq!(commission, None);
5037 }
5038
5039 #[rstest]
5040 fn test_resolve_inferred_fill_commission_propagates_failure() {
5041 let (order, report, instrument) = commission_fixtures();
5042 let client = CommissionStubClient::new(CommissionOutcome::Failure);
5043
5044 let result = ExecutionManager::resolve_inferred_fill_commission(
5045 Some(&client),
5046 &instrument,
5047 Quantity::from("5.0"),
5048 inferred_fill_price_and_liquidity(&order, &report, &instrument),
5049 );
5050
5051 assert!(result.is_err(), "a venue failure must not become Ok(None)");
5052 }
5053
5054 fn external_report_with_partial_fill(
5055 instrument: &InstrumentAny,
5056 ) -> (OrderStatusReport, FillReport) {
5057 let account_id = AccountId::from("STUB-001");
5058 let venue_order_id = VenueOrderId::from("V-EXT-1");
5059 let report = OrderStatusReport::new(
5060 account_id,
5061 instrument.id(),
5062 None,
5063 venue_order_id,
5064 OrderSide::Buy,
5065 OrderType::Limit,
5066 TimeInForce::Gtc,
5067 OrderStatus::Filled,
5068 Quantity::from("10.0"),
5069 Quantity::from("10.0"),
5070 UnixNanos::from(1),
5071 UnixNanos::from(1),
5072 UnixNanos::from(1),
5073 None,
5074 )
5075 .with_price(Price::from("100.00"))
5076 .with_avg_px(dec!(100.0));
5077
5078 let fill = FillReport::new(
5079 account_id,
5080 instrument.id(),
5081 venue_order_id,
5082 TradeId::from("T-EXT-1"),
5083 OrderSide::Buy,
5084 Quantity::from("4.0"),
5085 Price::from("100.00"),
5086 Money::new(0.1, Currency::USDT()),
5087 LiquiditySide::Taker,
5088 None,
5089 None,
5090 UnixNanos::from(1),
5091 UnixNanos::from(1),
5092 None,
5093 );
5094
5095 (report, fill)
5096 }
5097
5098 fn inferred_fills(events: &[OrderEventAny]) -> Vec<OrderFilled> {
5099 events
5100 .iter()
5101 .filter_map(|event| match event {
5102 OrderEventAny::Filled(filled) if filled.last_qty == Quantity::from("6.0") => {
5103 Some(filled.clone())
5104 }
5105 _ => None,
5106 })
5107 .collect()
5108 }
5109
5110 #[rstest]
5111 fn test_handle_external_order_applies_venue_commission_to_inferred_fill() {
5112 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
5113 let clock = Rc::new(RefCell::new(TestClock::new()));
5114 let cache = Rc::new(RefCell::new(Cache::default()));
5115 cache
5116 .borrow_mut()
5117 .add_instrument(instrument.clone())
5118 .expect("instrument is cacheable");
5119 let manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default());
5120 let (report, fill) = external_report_with_partial_fill(&instrument);
5121 let expected = Money::new(2.5, Currency::USDT());
5122 let client = CommissionStubClient::new(CommissionOutcome::Value(expected));
5123 let mut fill_queue = ReconciliationFillQueue::default();
5124
5125 let (events, _) = manager.handle_external_order(
5126 &report,
5127 AccountId::from("STUB-001"),
5128 &instrument,
5129 &[&fill],
5130 false,
5131 Some(&mut fill_queue),
5132 Some(&client),
5133 );
5134
5135 let inferred = inferred_fills(&events);
5136 assert_eq!(inferred.len(), 1, "one inferred fill covers the 6.0 gap");
5137 assert_eq!(inferred[0].commission, Some(expected));
5138 }
5139
5140 #[rstest]
5141 fn test_handle_external_order_skips_inferred_fill_when_commission_fails() {
5142 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
5143 let clock = Rc::new(RefCell::new(TestClock::new()));
5144 let cache = Rc::new(RefCell::new(Cache::default()));
5145 cache
5146 .borrow_mut()
5147 .add_instrument(instrument.clone())
5148 .expect("instrument is cacheable");
5149 let manager =
5150 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default());
5151 let (report, fill) = external_report_with_partial_fill(&instrument);
5152 let client = CommissionStubClient::new(CommissionOutcome::Failure);
5153 let mut fill_queue = ReconciliationFillQueue::default();
5154
5155 let (events, metadata) = manager.handle_external_order(
5156 &report,
5157 AccountId::from("STUB-001"),
5158 &instrument,
5159 &[&fill],
5160 false,
5161 Some(&mut fill_queue),
5162 Some(&client),
5163 );
5164
5165 assert!(events.is_empty());
5166 assert!(metadata.is_none());
5167 assert!(fill_queue.pending_fill_keys.is_empty());
5168 assert!(
5169 cache
5170 .borrow()
5171 .order(&ClientOrderId::from(report.venue_order_id.as_str()))
5172 .is_none(),
5173 "commission failure must precede external order cache mutation"
5174 );
5175
5176 let expected = Money::new(2.5, Currency::USDT());
5177 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
5178 let (retry_events, retry_metadata) = manager.handle_external_order(
5179 &report,
5180 AccountId::from("STUB-001"),
5181 &instrument,
5182 &[&fill],
5183 false,
5184 Some(&mut fill_queue),
5185 Some(&retry_client),
5186 );
5187 let inferred = inferred_fills(&retry_events);
5188
5189 assert!(retry_metadata.is_some());
5190 assert_eq!(inferred.len(), 1);
5191 assert_eq!(inferred[0].commission, Some(expected));
5192 assert_eq!(fill_queue.pending_fill_keys.len(), 1);
5193 assert!(
5194 cache
5195 .borrow()
5196 .order(&ClientOrderId::from(report.venue_order_id.as_str()))
5197 .is_some()
5198 );
5199 }
5200
5201 #[rstest]
5202 fn test_handle_external_order_without_explicit_fills_resolves_commission_before_cache() {
5203 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
5204 let clock = Rc::new(RefCell::new(TestClock::new()));
5205 let cache = Rc::new(RefCell::new(Cache::default()));
5206 cache
5207 .borrow_mut()
5208 .add_instrument(instrument.clone())
5209 .expect("instrument is cacheable");
5210 let manager =
5211 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default());
5212 let (report, _) = external_report_with_partial_fill(&instrument);
5213 let failing_client = CommissionStubClient::new(CommissionOutcome::Failure);
5214
5215 let (failed_events, failed_metadata) = manager.handle_external_order(
5216 &report,
5217 AccountId::from("STUB-001"),
5218 &instrument,
5219 &[],
5220 false,
5221 None,
5222 Some(&failing_client),
5223 );
5224
5225 assert!(failed_events.is_empty());
5226 assert!(failed_metadata.is_none());
5227 assert!(
5228 cache
5229 .borrow()
5230 .order(&ClientOrderId::from(report.venue_order_id.as_str()))
5231 .is_none()
5232 );
5233
5234 let expected = Money::new(4.0, Currency::USDT());
5235 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
5236 let (retry_events, retry_metadata) = manager.handle_external_order(
5237 &report,
5238 AccountId::from("STUB-001"),
5239 &instrument,
5240 &[],
5241 false,
5242 None,
5243 Some(&retry_client),
5244 );
5245 let fills: Vec<_> = retry_events
5246 .iter()
5247 .filter_map(|event| match event {
5248 OrderEventAny::Filled(fill) => Some(fill),
5249 _ => None,
5250 })
5251 .collect();
5252
5253 assert!(retry_metadata.is_some());
5254 assert_eq!(fills.len(), 1);
5255 assert_eq!(fills[0].last_qty, Quantity::from("10.0"));
5256 assert_eq!(fills[0].commission, Some(expected));
5257 }
5258
5259 #[rstest]
5260 fn test_handle_external_order_with_no_override_emits_fill_without_commission() {
5261 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
5262 let clock = Rc::new(RefCell::new(TestClock::new()));
5263 let cache = Rc::new(RefCell::new(Cache::default()));
5264 cache
5265 .borrow_mut()
5266 .add_instrument(instrument.clone())
5267 .expect("instrument is cacheable");
5268 let manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default());
5269 let (report, fill) = external_report_with_partial_fill(&instrument);
5270 let client = CommissionStubClient::new(CommissionOutcome::NoOverride);
5271 let mut fill_queue = ReconciliationFillQueue::default();
5272
5273 let (events, _) = manager.handle_external_order(
5274 &report,
5275 AccountId::from("STUB-001"),
5276 &instrument,
5277 &[&fill],
5278 false,
5279 Some(&mut fill_queue),
5280 Some(&client),
5281 );
5282
5283 let inferred = inferred_fills(&events);
5284 assert_eq!(inferred.len(), 1);
5285 assert_eq!(inferred[0].commission, None);
5286 }
5287
5288 #[rstest]
5289 fn test_cached_reconciliation_applies_explicit_fill_and_defers_failed_residual() {
5290 let (mut manager, _cache, order, mut report, instrument) = cached_commission_fixtures();
5291 report.avg_px = Some(dec!(60.0));
5292 let explicit_fill = FillReport::new(
5293 report.account_id,
5294 report.instrument_id,
5295 report.venue_order_id,
5296 TradeId::from("T-COMMISSION-EXPLICIT"),
5297 OrderSide::Buy,
5298 Quantity::from("4.0"),
5299 Price::from("100.0"),
5300 Money::new(0.25, Currency::USDT()),
5301 LiquiditySide::Taker,
5302 report.client_order_id,
5303 None,
5304 UnixNanos::from(1),
5305 UnixNanos::from(1),
5306 None,
5307 );
5308 let failing_client = CommissionStubClient::new(CommissionOutcome::Failure);
5309 let mut fill_queue = ReconciliationFillQueue::default();
5310
5311 let first_events = manager.reconcile_order_with_fills(
5312 &order,
5313 &report,
5314 &[&explicit_fill],
5315 Some(&instrument),
5316 &mut fill_queue,
5317 Some(&failing_client),
5318 );
5319 let mut working = order;
5320 for event in &first_events {
5321 working
5322 .apply(event.clone())
5323 .expect("explicit fill projects cleanly");
5324 }
5325
5326 assert_eq!(first_events.len(), 1);
5327 let OrderEventAny::Filled(explicit) = &first_events[0] else {
5328 panic!("expected the valid explicit fill");
5329 };
5330 assert_eq!(explicit.last_qty, Quantity::from("4.0"));
5331 assert_eq!(
5332 explicit.commission,
5333 Some(Money::new(0.25, Currency::USDT()))
5334 );
5335 assert_eq!(working.status(), OrderStatus::PartiallyFilled);
5336
5337 let expected = Money::new(1.5, Currency::USDT());
5338 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
5339 let retry_events = manager.reconcile_order_with_fills(
5340 &working,
5341 &report,
5342 &[],
5343 Some(&instrument),
5344 &mut fill_queue,
5345 Some(&retry_client),
5346 );
5347
5348 assert_eq!(retry_events.len(), 1);
5349 let OrderEventAny::Filled(residual) = &retry_events[0] else {
5350 panic!("expected the inferred residual fill");
5351 };
5352 assert_eq!(residual.last_qty, Quantity::from("6.0"));
5353 assert_eq!(residual.last_px, Price::from("33.33"));
5354 assert_eq!(residual.commission, Some(expected));
5355 assert_eq!(
5356 *retry_client.seen.borrow(),
5357 Some((
5358 Quantity::from("6.0"),
5359 residual.last_px,
5360 residual.liquidity_side,
5361 )),
5362 "commission must use the exact price and liquidity carried by the residual fill"
5363 );
5364 }
5365
5366 #[rstest]
5367 fn test_cached_snapshot_without_instrument_preserves_terminal_transition() {
5368 let (mut manager, _cache, order, mut report, _instrument) = cached_commission_fixtures();
5369 report.order_status = OrderStatus::Canceled;
5370 let explicit_fill = FillReport::new(
5371 report.account_id,
5372 report.instrument_id,
5373 report.venue_order_id,
5374 TradeId::from("T-MISSING-INSTRUMENT"),
5375 OrderSide::Buy,
5376 Quantity::from("4.0"),
5377 Price::from("100.0"),
5378 Money::new(0.25, Currency::USDT()),
5379 LiquiditySide::Taker,
5380 report.client_order_id,
5381 None,
5382 UnixNanos::from(1),
5383 UnixNanos::from(1),
5384 None,
5385 );
5386 let mut fill_queue = ReconciliationFillQueue::default();
5387
5388 let events = manager.reconcile_order_with_fills(
5389 &order,
5390 &report,
5391 &[&explicit_fill],
5392 None,
5393 &mut fill_queue,
5394 None,
5395 );
5396
5397 assert_eq!(events.len(), 1);
5398 assert!(matches!(events[0], OrderEventAny::Canceled(_)));
5399 }
5400
5401 #[rstest]
5402 fn test_continuous_reconciliation_uses_source_client_and_retries_commission() {
5403 let (mut manager, _cache, order, report, _instrument) = cached_commission_fixtures();
5404 let client_id = ClientId::from("STUB");
5405 let check = OpenOrderReportCheck {
5406 command: GenerateOrderStatusReports::new(
5407 UUID4::new(),
5408 UnixNanos::from(1),
5409 true,
5410 None,
5411 None,
5412 None,
5413 None,
5414 None,
5415 ),
5416 filtered_orders: vec![order],
5417 client_coverage: IndexMap::from([(
5418 report.client_order_id.unwrap(),
5419 ReportClientCoverage::Resolved(IndexSet::from([client_id])),
5420 )]),
5421 start: None,
5422 };
5423 let queried_clients = IndexSet::from([client_id]);
5424 let failed_clients = IndexSet::new();
5425 let failing_client = CommissionStubClient::new(CommissionOutcome::Failure);
5426
5427 let failed = manager.reconcile_open_order_reports(
5428 &check,
5429 vec![SourcedOrderStatusReport {
5430 client_id,
5431 report: report.clone(),
5432 }],
5433 &queried_clients,
5434 &failed_clients,
5435 &[&failing_client],
5436 );
5437
5438 assert!(failed.events.is_empty());
5439
5440 let expected = Money::new(1.5, Currency::USDT());
5441 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
5442 let retry = manager.reconcile_open_order_reports(
5443 &check,
5444 vec![SourcedOrderStatusReport { client_id, report }],
5445 &queried_clients,
5446 &failed_clients,
5447 &[&retry_client],
5448 );
5449
5450 assert_eq!(retry.events.len(), 1);
5451 let OrderEventAny::Filled(fill) = &retry.events[0] else {
5452 panic!("expected inferred fill on valid retry");
5453 };
5454 assert_eq!(fill.last_qty, Quantity::from("10.0"));
5455 assert_eq!(fill.commission, Some(expected));
5456 }
5457
5458 #[rstest]
5459 fn test_targeted_reconciliation_uses_source_client_and_retries_commission() {
5460 let (mut manager, _cache, _order, report, _instrument) = cached_commission_fixtures();
5461 let client_order_id = report.client_order_id.unwrap();
5462 let client_id = ClientId::from("STUB");
5463 let failing_client = CommissionStubClient::new(CommissionOutcome::Failure);
5464
5465 let failed = manager.reconcile_targeted_order_reports(
5466 vec![TargetedOrderReportResult {
5467 client_order_id,
5468 client_id: Some(client_id),
5469 report: Some(report.clone()),
5470 coverage_complete: true,
5471 }],
5472 &[&failing_client],
5473 );
5474
5475 assert!(failed.is_empty());
5476
5477 let expected = Money::new(1.5, Currency::USDT());
5478 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
5479 let retry = manager.reconcile_targeted_order_reports(
5480 vec![TargetedOrderReportResult {
5481 client_order_id,
5482 client_id: Some(client_id),
5483 report: Some(report),
5484 coverage_complete: true,
5485 }],
5486 &[&retry_client],
5487 );
5488
5489 assert_eq!(retry.len(), 1);
5490 let OrderEventAny::Filled(fill) = &retry[0] else {
5491 panic!("expected inferred fill on valid targeted retry");
5492 };
5493 assert_eq!(fill.last_qty, Quantity::from("10.0"));
5494 assert_eq!(fill.commission, Some(expected));
5495 }
5496
5497 #[rstest]
5498 fn test_clear_recon_tracking_removes_targeted_query() {
5499 let clock = Rc::new(RefCell::new(TestClock::new()));
5500 let cache = Rc::new(RefCell::new(Cache::default()));
5501 let mut manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default());
5502 let client_order_id = ClientOrderId::from("O-TARGETED-CLEAR");
5503 manager.targeted_order_queries.insert(client_order_id);
5504
5505 manager.clear_recon_tracking(&client_order_id, true);
5506
5507 assert!(manager.targeted_order_queries.is_empty());
5508 }
5509
5510 #[rstest]
5511 fn test_register_inflight_skips_filtered_order() {
5512 let client_order_id = ClientOrderId::from("O-FILTERED-REGISTER");
5513 let clock = Rc::new(RefCell::new(TestClock::new()));
5514 let cache = Rc::new(RefCell::new(Cache::default()));
5515 let mut manager = ExecutionManager::new(
5516 clock,
5517 cache,
5518 ExecutionManagerConfig {
5519 filtered_client_order_ids: IndexSet::from([client_order_id]),
5520 ..Default::default()
5521 },
5522 );
5523
5524 manager.register_inflight(client_order_id);
5525
5526 assert!(!manager.inflight_checks.contains_key(&client_order_id));
5527 assert!(!manager.recon_check_retries.contains_key(&client_order_id));
5528 }
5529
5530 #[rstest]
5531 #[cfg_attr(
5532 not(all(feature = "simulation", madsim)),
5533 tokio::test(start_paused = true)
5534 )]
5535 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
5536 async fn test_inflight_check_retires_order_filtered_after_registration() {
5537 let client_order_id = ClientOrderId::from("O-FILTERED-LATE");
5538 let clock = Rc::new(RefCell::new(TestClock::new()));
5539 let cache = Rc::new(RefCell::new(Cache::default()));
5540 let mut manager = ExecutionManager::new(
5541 clock,
5542 cache,
5543 ExecutionManagerConfig {
5544 inflight_threshold_ms: 100,
5545 ..Default::default()
5546 },
5547 );
5548 manager.register_inflight(client_order_id);
5549 manager
5550 .config
5551 .filtered_client_order_ids
5552 .insert(client_order_id);
5553 dst::time::sleep(Duration::from_millis(101)).await;
5554
5555 let first = manager.check_inflight_orders();
5556
5557 assert!(first.events.is_empty());
5558 assert!(first.queries.is_empty());
5559 assert!(!manager.inflight_checks.contains_key(&client_order_id));
5560 assert!(!manager.recon_check_retries.contains_key(&client_order_id));
5561
5562 dst::time::sleep(Duration::from_millis(101)).await;
5563 let second = manager.check_inflight_orders();
5564 assert!(second.events.is_empty());
5565 assert!(second.queries.is_empty());
5566 assert!(!manager.inflight_checks.contains_key(&client_order_id));
5567 }
5568
5569 #[rstest]
5570 #[case(false, OrderStatus::PendingUpdate, true, true, true)]
5571 #[case(false, OrderStatus::Accepted, false, true, true)]
5572 #[case(false, OrderStatus::Canceled, false, true, false)]
5573 #[case(true, OrderStatus::PendingCancel, true, true, true)]
5574 #[case(true, OrderStatus::Accepted, false, true, true)]
5575 #[case(true, OrderStatus::Filled, false, true, false)]
5576 fn test_observe_order_status_report_tracking_matrix(
5577 #[case] with_fills: bool,
5578 #[case] status: OrderStatus,
5579 #[case] expect_inflight: bool,
5580 #[case] expect_activity: bool,
5581 #[case] expect_last_query: bool,
5582 ) {
5583 let client_order_id = ClientOrderId::from("O-STATUS-MATRIX");
5584 let clock = Rc::new(RefCell::new(TestClock::new()));
5585 let cache = Rc::new(RefCell::new(Cache::default()));
5586 let mut manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default());
5587 manager.register_inflight(client_order_id);
5588 manager.order_query_recency.mark(client_order_id);
5589 manager
5590 .missing_order_coverage_warnings
5591 .insert(client_order_id);
5592 manager.unresolved_order_coverage.insert(client_order_id);
5593 manager.targeted_order_queries.insert(client_order_id);
5594 let order_report = OrderStatusReport::new(
5595 AccountId::from("TEST-001"),
5596 crypto_perpetual_ethusdt().id(),
5597 Some(client_order_id),
5598 VenueOrderId::from("V-STATUS-MATRIX"),
5599 OrderSide::Buy,
5600 OrderType::Limit,
5601 TimeInForce::Gtc,
5602 status,
5603 Quantity::from("10.0"),
5604 Quantity::from("0.0"),
5605 UnixNanos::from(1_000),
5606 UnixNanos::from(1_000),
5607 UnixNanos::from(1_000),
5608 None,
5609 );
5610 let report = if with_fills {
5611 ExecutionReport::OrderWithFills(Box::new(order_report), Vec::new())
5612 } else {
5613 ExecutionReport::Order(Box::new(order_report))
5614 };
5615
5616 manager.observe_execution_report(&report);
5617
5618 assert_eq!(
5619 manager.inflight_checks.contains_key(&client_order_id),
5620 expect_inflight,
5621 );
5622 assert_eq!(
5623 manager.recon_check_retries.contains_key(&client_order_id),
5624 expect_inflight,
5625 );
5626 assert_eq!(
5627 manager.order_local_activity.contains_key(&client_order_id),
5628 expect_activity,
5629 );
5630 assert_eq!(
5631 manager.order_query_recency.contains_key(&client_order_id),
5632 expect_last_query,
5633 );
5634 assert_eq!(
5635 manager
5636 .missing_order_coverage_warnings
5637 .contains(&client_order_id),
5638 expect_inflight,
5639 );
5640 assert_eq!(
5641 manager.unresolved_order_coverage.contains(&client_order_id),
5642 expect_inflight,
5643 );
5644 assert_eq!(
5645 manager.targeted_order_queries.contains(&client_order_id),
5646 expect_inflight,
5647 );
5648 }
5649
5650 #[rstest]
5651 fn test_superseded_cancel_report_preserves_missing_order_grace() {
5652 let client_order_id = ClientOrderId::from("O-CANCEL-REPLACE");
5653 let old_venue_order_id = VenueOrderId::from("V-CANCEL-REPLACE-OLD");
5654 let new_venue_order_id = VenueOrderId::from("V-CANCEL-REPLACE-NEW");
5655 let account_id = AccountId::from("TEST-001");
5656 let client_id = ClientId::from("TEST");
5657 let instrument_id = crypto_perpetual_ethusdt().id();
5658 let clock = Rc::new(RefCell::new(TestClock::new()));
5659 let cache = Rc::new(RefCell::new(Cache::default()));
5660 insert_accepted_limit_order(
5661 &cache,
5662 client_order_id,
5663 old_venue_order_id,
5664 instrument_id,
5665 client_id,
5666 );
5667
5668 let order = cache.borrow().order_owned(&client_order_id).unwrap();
5669 let pending_update = OrderPendingUpdateSpec::builder()
5670 .trader_id(order.trader_id())
5671 .strategy_id(order.strategy_id())
5672 .instrument_id(order.instrument_id())
5673 .client_order_id(client_order_id)
5674 .account_id(account_id)
5675 .venue_order_id(old_venue_order_id)
5676 .build();
5677 cache
5678 .borrow_mut()
5679 .update_order(&OrderEventAny::PendingUpdate(pending_update))
5680 .unwrap();
5681 let order = cache.borrow().order_owned(&client_order_id).unwrap();
5682 let updated = OrderUpdatedSpec::builder()
5683 .trader_id(order.trader_id())
5684 .strategy_id(order.strategy_id())
5685 .instrument_id(order.instrument_id())
5686 .client_order_id(client_order_id)
5687 .quantity(order.quantity())
5688 .venue_order_id(new_venue_order_id)
5689 .account_id(account_id)
5690 .build();
5691 cache
5692 .borrow_mut()
5693 .update_order(&OrderEventAny::Updated(updated))
5694 .unwrap();
5695
5696 let mut manager = ExecutionManager::new(
5697 clock,
5698 cache.clone(),
5699 ExecutionManagerConfig {
5700 open_check_missing_retries: 1,
5701 ..Default::default()
5702 },
5703 );
5704 manager.record_local_activity(client_order_id);
5705 assert!(
5706 manager
5707 .prepare_missing_order_query(client_order_id)
5708 .is_none()
5709 );
5710
5711 let report = OrderStatusReport::new(
5712 account_id,
5713 instrument_id,
5714 Some(client_order_id),
5715 old_venue_order_id,
5716 OrderSide::Buy,
5717 OrderType::Limit,
5718 TimeInForce::Gtc,
5719 OrderStatus::Canceled,
5720 Quantity::from("10.0"),
5721 Quantity::from("0.0"),
5722 UnixNanos::from(1_000),
5723 UnixNanos::from(2_000),
5724 UnixNanos::from(3_000),
5725 None,
5726 );
5727
5728 manager.observe_execution_report(&ExecutionReport::Order(Box::new(report.clone())));
5729 let order = cache.borrow().order_owned(&client_order_id).unwrap();
5730 let events =
5731 generate_reconciliation_order_events(&order, &report, None, UnixNanos::from(1_000));
5732
5733 assert!(events.is_empty());
5734 assert_eq!(order.status(), OrderStatus::Accepted);
5735 assert_eq!(order.venue_order_id(), Some(new_venue_order_id));
5736 assert!(manager.order_local_activity.contains_key(&client_order_id));
5737 assert!(
5738 manager
5739 .prepare_missing_order_query(client_order_id)
5740 .is_none()
5741 );
5742 assert_eq!(manager.recon_check_retry_count(&client_order_id), 0);
5743 }
5744
5745 #[rstest]
5746 #[cfg_attr(
5747 not(all(feature = "simulation", madsim)),
5748 tokio::test(start_paused = true)
5749 )]
5750 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
5751 async fn test_prune_order_local_activity_uses_open_check_threshold() {
5752 let old_id = ClientOrderId::from("O-ACTIVITY-OLD");
5753 let fresh_id = ClientOrderId::from("O-ACTIVITY-FRESH");
5754 let clock = Rc::new(RefCell::new(TestClock::new()));
5755 let cache = Rc::new(RefCell::new(Cache::default()));
5756 let mut manager = ExecutionManager::new(
5757 clock,
5758 cache,
5759 ExecutionManagerConfig {
5760 open_check_threshold_ns: 100_000_000,
5761 ..Default::default()
5762 },
5763 );
5764 manager.record_local_activity(old_id);
5765 dst::time::sleep(Duration::from_millis(101)).await;
5766 manager.record_local_activity(fresh_id);
5767
5768 manager.prune_order_local_activity();
5769
5770 assert!(!manager.order_local_activity.contains_key(&old_id));
5771 assert!(manager.order_local_activity.contains_key(&fresh_id));
5772 }
5773
5774 #[rstest]
5775 fn test_prepare_open_order_report_check_builds_bulk_command_with_config() {
5776 let lookback_mins = 5_u64;
5777 let lookback_ns = lookback_mins * 60 * NANOSECONDS_IN_SECOND;
5778 let clock = Rc::new(RefCell::new(TestClock::new()));
5779 let cache = Rc::new(RefCell::new(Cache::default()));
5780 let mut manager = ExecutionManager::new(
5781 clock.clone(),
5782 cache.clone(),
5783 ExecutionManagerConfig {
5784 open_check_lookback_mins: Some(lookback_mins),
5785 open_check_open_only: false,
5786 reconciliation_instrument_ids: IndexSet::from([crypto_perpetual_ethusdt().id()]),
5787 ..Default::default()
5788 },
5789 );
5790 let included_id = ClientOrderId::from("O-REPORT-001");
5791 let excluded_id = ClientOrderId::from("O-REPORT-002");
5792 let included_instrument_id = crypto_perpetual_ethusdt().id();
5793 let excluded_instrument_id = xbtusd_bitmex().id();
5794
5795 cache
5796 .borrow_mut()
5797 .add_instrument(InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt()))
5798 .unwrap();
5799 cache
5800 .borrow_mut()
5801 .add_instrument(InstrumentAny::CryptoPerpetual(xbtusd_bitmex()))
5802 .unwrap();
5803 insert_accepted_limit_order(
5804 &cache,
5805 included_id,
5806 VenueOrderId::from("V-REPORT-001"),
5807 included_instrument_id,
5808 ClientId::from("BINANCE"),
5809 );
5810 insert_accepted_limit_order(
5811 &cache,
5812 excluded_id,
5813 VenueOrderId::from("V-REPORT-002"),
5814 excluded_instrument_id,
5815 ClientId::from("BITMEX"),
5816 );
5817 clock
5818 .borrow_mut()
5819 .advance_time(UnixNanos::from(lookback_ns * 2), true);
5820
5821 let ts_now = clock.borrow().timestamp_ns();
5822 let command_id = UUID4::new();
5823 let check = manager.prepare_open_order_report_check(command_id, &[]);
5824
5825 assert_eq!(check.command.command_id, command_id);
5826 assert_eq!(check.command.ts_init, ts_now);
5827 assert!(!check.command.open_only);
5828 assert_eq!(check.command.instrument_id, None);
5829 assert_eq!(
5830 check.command.start,
5831 Some(ts_now.saturating_sub_ns(lookback_ns))
5832 );
5833 assert_eq!(check.command.end, None);
5834 assert_eq!(check.command.log_receipt_level, LogLevel::Debug);
5835 assert_eq!(check.start, check.command.start);
5836 assert_eq!(check.filtered_orders.len(), 1);
5837 assert_eq!(check.filtered_orders[0].client_order_id(), included_id);
5838 }
5839
5840 #[rstest]
5841 fn test_prepare_position_report_check_builds_bulk_command_with_coverage() {
5842 let clock = Rc::new(RefCell::new(TestClock::new()));
5843 let cache = Rc::new(RefCell::new(Cache::default()));
5844 let manager = ExecutionManager::new(
5845 clock.clone(),
5846 cache.clone(),
5847 ExecutionManagerConfig {
5848 reconciliation_instrument_ids: IndexSet::from([crypto_perpetual_ethusdt().id()]),
5849 ..Default::default()
5850 },
5851 );
5852 let included_instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
5853 let excluded_instrument = InstrumentAny::CryptoPerpetual(xbtusd_bitmex());
5854
5855 cache
5856 .borrow_mut()
5857 .add_instrument(included_instrument.clone())
5858 .unwrap();
5859 cache
5860 .borrow_mut()
5861 .add_instrument(excluded_instrument.clone())
5862 .unwrap();
5863 let included_position = insert_open_position(
5864 &cache,
5865 &included_instrument,
5866 PositionId::from("P-REPORT-001"),
5867 OrderSide::Buy,
5868 "5.0",
5869 "3000.00",
5870 );
5871 insert_open_position(
5872 &cache,
5873 &excluded_instrument,
5874 PositionId::from("P-REPORT-002"),
5875 OrderSide::Buy,
5876 "2.0",
5877 "40000.00",
5878 );
5879
5880 let ts_now = clock.borrow().timestamp_ns();
5881 let command_id = UUID4::new();
5882 let check = manager.prepare_position_report_check(command_id, &[]);
5883 let key = (
5884 included_position.instrument_id,
5885 included_position.account_id,
5886 );
5887
5888 assert_eq!(check.command.command_id, command_id);
5889 assert_eq!(check.command.ts_init, ts_now);
5890 assert_eq!(check.command.instrument_id, None);
5891 assert_eq!(check.command.start, None);
5892 assert_eq!(check.command.end, None);
5893 assert_eq!(check.command.log_receipt_level, LogLevel::Debug);
5894 assert_eq!(check.client_coverage.len(), 1);
5895 assert!(check.client_coverage.contains_key(&key));
5896 assert_eq!(check.activity_revisions.get(&key), Some(&0));
5897 }
5898
5899 #[rstest]
5900 #[cfg_attr(
5901 not(all(feature = "simulation", madsim)),
5902 tokio::test(start_paused = true)
5903 )]
5904 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
5905 async fn test_position_report_check_defers_activity_recorded_during_delayed_request() {
5906 let clock = Rc::new(RefCell::new(TestClock::new()));
5907 let cache = Rc::new(RefCell::new(Cache::default()));
5908 let mut manager = ExecutionManager::new(
5909 clock,
5910 cache.clone(),
5911 ExecutionManagerConfig {
5912 position_check_threshold_ns: 5_000_000_000,
5913 ..Default::default()
5914 },
5915 );
5916 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
5917 let instrument_id = instrument.id();
5918 let position = insert_open_position(
5919 &cache,
5920 &instrument,
5921 PositionId::from("P-ACTIVITY-DURING-REQUEST"),
5922 OrderSide::Buy,
5923 "5.0",
5924 "3000.00",
5925 );
5926 cache
5927 .borrow_mut()
5928 .add_instrument(instrument.clone())
5929 .unwrap();
5930 let account_id = position.account_id;
5931 let check = manager.prepare_position_report_check(UUID4::new(), &[]);
5932 let report = PositionStatusReport::new(
5933 account_id,
5934 instrument_id,
5935 PositionSideSpecified::Long,
5936 Quantity::from("5.0"),
5937 UnixNanos::from(1_000_000),
5938 UnixNanos::from(1_000_000),
5939 None,
5940 None,
5941 Some(Decimal::from(3000)),
5942 );
5943
5944 let closed_position = close_long_position(
5945 position,
5946 &instrument,
5947 TradeId::from("T-ACTIVITY-DURING-REQUEST"),
5948 );
5949 cache
5950 .borrow_mut()
5951 .update_position(&closed_position)
5952 .unwrap();
5953 manager.record_position_activity(instrument_id, account_id);
5954
5955 dst::time::sleep(Duration::from_secs(6)).await;
5957
5958 let events = manager.reconcile_position_reports(
5959 &check,
5960 vec![report],
5961 &IndexSet::new(),
5962 &IndexSet::new(),
5963 );
5964
5965 assert!(
5966 !events.iter().any(|event| {
5967 matches!(
5968 event,
5969 OrderEventAny::Filled(fill)
5970 if fill.order_side == OrderSide::Buy
5971 && fill.last_qty == Quantity::from("5.0")
5972 )
5973 }),
5974 "activity recorded after the request started must defer A's stale report",
5975 );
5976 }
5977
5978 #[rstest]
5979 fn test_position_report_check_does_not_defer_activity_recorded_before_request() {
5980 let clock = Rc::new(RefCell::new(TestClock::new()));
5981 let cache = Rc::new(RefCell::new(Cache::default()));
5982 let mut manager = ExecutionManager::new(
5983 clock,
5984 cache.clone(),
5985 ExecutionManagerConfig {
5986 position_check_threshold_ns: 0,
5987 ..Default::default()
5988 },
5989 );
5990 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
5991 let instrument_id = instrument.id();
5992 let position = insert_open_position(
5993 &cache,
5994 &instrument,
5995 PositionId::from("P-ACTIVITY-BEFORE-REQUEST"),
5996 OrderSide::Buy,
5997 "5.0",
5998 "3000.00",
5999 );
6000 cache.borrow_mut().add_instrument(instrument).unwrap();
6001 let account_id = position.account_id;
6002 manager.record_position_activity(instrument_id, account_id);
6003 let check = manager.prepare_position_report_check(UUID4::new(), &[]);
6004 let report = PositionStatusReport::new(
6005 account_id,
6006 instrument_id,
6007 PositionSideSpecified::Long,
6008 Quantity::from("10.0"),
6009 UnixNanos::from(1_000_000),
6010 UnixNanos::from(1_000_000),
6011 None,
6012 None,
6013 Some(Decimal::from(3000)),
6014 );
6015
6016 let events = manager.reconcile_position_reports(
6017 &check,
6018 vec![report],
6019 &IndexSet::new(),
6020 &IndexSet::new(),
6021 );
6022
6023 let fills: Vec<_> = events
6024 .iter()
6025 .filter_map(|event| match event {
6026 OrderEventAny::Filled(fill) => Some(fill),
6027 _ => None,
6028 })
6029 .collect();
6030
6031 assert_eq!(fills.len(), 1);
6032 assert_eq!(fills[0].order_side, OrderSide::Buy);
6033 assert_eq!(fills[0].last_qty, Quantity::from("5.0"));
6034 assert_eq!(fills[0].commission, None);
6035 }
6036
6037 #[rstest]
6038 fn test_mass_status_projects_companion_fill_before_void_correction() {
6039 let clock = Rc::new(RefCell::new(TestClock::new()));
6040 let cache = Rc::new(RefCell::new(Cache::default()));
6041 let mut manager =
6042 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default());
6043 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
6044 let client_order_id = ClientOrderId::from("O-MASS-VOID-001");
6045 let venue_order_id = VenueOrderId::from("V-MASS-VOID-001");
6046 let account_id = AccountId::from("TEST-001");
6047 cache
6048 .borrow_mut()
6049 .add_instrument(instrument.clone())
6050 .unwrap();
6051 insert_accepted_limit_order(
6052 &cache,
6053 client_order_id,
6054 venue_order_id,
6055 instrument.id(),
6056 ClientId::from("BINANCE"),
6057 );
6058 let order = cache.borrow().order_owned(&client_order_id).unwrap();
6059 let initial_fill = TestOrderEventStubs::filled(
6060 &order,
6061 &instrument,
6062 Some(TradeId::from("T-MASS-VOID-INITIAL")),
6063 None,
6064 Some(Price::from("100.0")),
6065 Some(Quantity::from("6.0")),
6066 Some(LiquiditySide::Taker),
6067 None,
6068 None,
6069 Some(account_id),
6070 );
6071 cache.borrow_mut().update_order(&initial_fill).unwrap();
6072 let order = cache.borrow().order_owned(&client_order_id).unwrap();
6073 let report = OrderStatusReport::new(
6074 account_id,
6075 instrument.id(),
6076 Some(client_order_id),
6077 venue_order_id,
6078 OrderSide::Buy,
6079 OrderType::Limit,
6080 TimeInForce::Gtc,
6081 OrderStatus::Canceled,
6082 Quantity::from("10.0"),
6083 Quantity::from("5.0"),
6084 UnixNanos::from(1_000),
6085 UnixNanos::from(1_000),
6086 UnixNanos::from(1_000),
6087 None,
6088 );
6089 let companion_fill = FillReport::new(
6090 account_id,
6091 instrument.id(),
6092 venue_order_id,
6093 TradeId::from("T-MASS-VOID-COMPANION"),
6094 OrderSide::Buy,
6095 Quantity::from("1.0"),
6096 Price::from("100.0"),
6097 Money::zero(instrument.quote_currency()),
6098 LiquiditySide::Taker,
6099 Some(client_order_id),
6100 None,
6101 UnixNanos::from(900),
6102 UnixNanos::from(1_000),
6103 None,
6104 );
6105
6106 let mut fill_queue = ReconciliationFillQueue::default();
6107 let events = manager.reconcile_order_with_fills(
6108 &order,
6109 &report,
6110 &[&companion_fill],
6111 Some(&instrument),
6112 &mut fill_queue,
6113 None,
6114 );
6115 let mut projected = order;
6116 for event in &events {
6117 projected.apply(event.clone()).unwrap();
6118 }
6119
6120 assert!(matches!(events[0], OrderEventAny::Filled(_)));
6121 assert_eq!(
6122 events
6123 .iter()
6124 .filter(|event| matches!(event, OrderEventAny::FillVoided(_)))
6125 .count(),
6126 2
6127 );
6128 assert_eq!(projected.status(), OrderStatus::Canceled);
6129 assert_eq!(projected.filled_qty(), Quantity::from("5.0"));
6130 assert_eq!(projected.voided_qty(), Quantity::from("2.0"));
6131 }
6132
6133 fn insert_accepted_limit_order(
6134 cache: &Rc<RefCell<Cache>>,
6135 client_order_id: ClientOrderId,
6136 venue_order_id: VenueOrderId,
6137 instrument_id: InstrumentId,
6138 client_id: ClientId,
6139 ) {
6140 let account_id = AccountId::from("TEST-001");
6141 let order = OrderTestBuilder::new(OrderType::Limit)
6142 .client_order_id(client_order_id)
6143 .instrument_id(instrument_id)
6144 .quantity(Quantity::from("10.0"))
6145 .price(Price::from("100.0"))
6146 .build();
6147 let submitted = TestOrderEventStubs::submitted(&order, account_id);
6148 cache
6149 .borrow_mut()
6150 .add_order(order, None, Some(client_id), false)
6151 .unwrap();
6152 let order = cache.borrow_mut().update_order(&submitted).unwrap();
6153 let accepted = TestOrderEventStubs::accepted(&order, account_id, venue_order_id);
6154 cache.borrow_mut().update_order(&accepted).unwrap();
6155 }
6156
6157 fn insert_open_position(
6158 cache: &Rc<RefCell<Cache>>,
6159 instrument: &InstrumentAny,
6160 position_id: PositionId,
6161 side: OrderSide,
6162 quantity: &str,
6163 price: &str,
6164 ) -> Position {
6165 let order = OrderTestBuilder::new(OrderType::Market)
6166 .instrument_id(instrument.id())
6167 .side(side)
6168 .quantity(Quantity::from(quantity))
6169 .build();
6170 let fill = TestOrderEventStubs::filled(
6171 &order,
6172 instrument,
6173 Some(TradeId::new("T-REPORT-001")),
6174 Some(position_id),
6175 Some(Price::from(price)),
6176 Some(Quantity::from(quantity)),
6177 None,
6178 None,
6179 None,
6180 Some(AccountId::from("TEST-001")),
6181 );
6182 let order_filled: OrderFilled = fill.into();
6183 let position = Position::new(instrument, order_filled);
6184 cache
6185 .borrow_mut()
6186 .add_position(&position, OmsType::Hedging)
6187 .unwrap();
6188 position
6189 }
6190
6191 fn close_long_position(
6192 mut position: Position,
6193 instrument: &InstrumentAny,
6194 trade_id: TradeId,
6195 ) -> Position {
6196 let order = OrderTestBuilder::new(OrderType::Market)
6197 .instrument_id(instrument.id())
6198 .side(OrderSide::Sell)
6199 .quantity(position.quantity)
6200 .build();
6201 let fill = TestOrderEventStubs::filled(
6202 &order,
6203 instrument,
6204 Some(trade_id),
6205 Some(position.id),
6206 Some(Price::from("3000.00")),
6207 Some(position.quantity),
6208 None,
6209 None,
6210 None,
6211 Some(position.account_id),
6212 );
6213 let order_filled: OrderFilled = fill.into();
6214 position.apply(&order_filled);
6215 position
6216 }
6217}