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};
26#[cfg(feature = "node")]
27use nautilus_common::messages::execution::report::GenerateFillReports;
28use nautilus_common::{
29 cache::Cache,
30 clients::{DEFAULT_POSITION_RECONCILIATION_TOLERANCE, ExecutionClient},
31 clock::Clock,
32 config::{ConfigError, ConfigErrorCollector, ConfigResult},
33 enums::{LogColor, LogLevel},
34 live::dst,
35 log_info,
36 messages::{
37 ExecutionReport,
38 execution::{
39 QueryOrder, TradingCommand,
40 report::{
41 GenerateOrderStatusReport, GenerateOrderStatusReports,
42 GeneratePositionStatusReports,
43 },
44 },
45 },
46 msgbus::{self, MessagingSwitchboard, switchboard},
47};
48use nautilus_core::{
49 UUID4, UnixNanos,
50 datetime::{checked_mins_to_nanos, checked_mins_to_secs, mins_to_nanos, mins_to_secs},
51};
52#[cfg(feature = "node")]
53use nautilus_execution::reconciliation::create_inferred_reconciliation_trade_id;
54use nautilus_execution::{
55 engine::ExecutionEngine,
56 reconciliation::{
57 calculate_reconciliation_price, create_inferred_fill_for_qty,
58 create_position_reconciliation_venue_order_id, create_reconciliation_rejected,
59 create_reconciliation_triggered, generate_external_order_status_events_with_commission,
60 generate_reconciliation_order_pre_fill_events,
61 generate_reconciliation_order_snapshot_events_with_commission,
62 incremental_inferred_fill_price_and_liquidity, inferred_fill_price_and_liquidity,
63 process_mass_status_for_reconciliation,
64 process_mass_status_for_reconciliation_without_synthetic_reports,
65 reconcile_order_report_with_commission, should_reconciliation_update,
66 },
67};
68#[cfg(feature = "node")]
69use nautilus_model::position::PositionReplayEvent;
70#[cfg(feature = "node")]
71use nautilus_model::types::{money::MoneyRaw, quantity::QuantityRaw};
72use nautilus_model::{
73 enums::{LiquiditySide, OmsType, OrderSide, OrderStatus, OrderType, TimeInForce},
74 events::{OrderCanceled, OrderEventAny, OrderFilled, OrderInitialized},
75 identifiers::{
76 AccountId, ClientId, ClientOrderId, InstrumentId, PositionId, StrategyId, TradeId,
77 TraderId, VenueOrderId,
78 },
79 instruments::{Instrument, InstrumentAny},
80 orders::{Order, OrderAny, TRIGGERABLE_ORDER_TYPES},
81 position::Position,
82 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
83 types::{Money, Price, Quantity},
84};
85use rust_decimal::Decimal;
86use ustr::Ustr;
87
88use super::recency::RecencyMap;
89
90static TAG_VENUE: LazyLock<Ustr> = LazyLock::new(|| Ustr::from("VENUE"));
92
93static TAG_RECONCILIATION: LazyLock<Ustr> = LazyLock::new(|| Ustr::from("RECONCILIATION"));
95
96pub type InstrumentAccountKey = (InstrumentId, AccountId);
102type AccountInstrumentKey = (AccountId, InstrumentId);
103type AccountInstrumentStrategyKey = (AccountId, InstrumentId, StrategyId);
104type FillKey = (AccountId, InstrumentId, TradeId);
105
106#[expect(clippy::too_many_arguments)]
107fn build_cross_zero_leg_report(
108 instrument: &InstrumentAny,
109 account_id: AccountId,
110 instrument_id: InstrumentId,
111 order_side: OrderSide,
112 quantity: Decimal,
113 avg_px: Decimal,
114 venue_position_id: Option<PositionId>,
115 tag: &str,
116 ts_now: UnixNanos,
117 venue_ts_last: UnixNanos,
118) -> Option<OrderStatusReport> {
119 let order_qty = Quantity::from_decimal_dp(quantity, instrument.size_precision()).ok()?;
120 let fill_price = Price::from_decimal_dp(avg_px, instrument.price_precision()).ok();
121 let venue_order_id = create_position_reconciliation_venue_order_id(
122 account_id,
123 instrument_id,
124 order_side,
125 OrderType::Market,
126 order_qty,
127 fill_price,
128 venue_position_id,
129 Some(tag),
130 venue_ts_last,
131 );
132
133 let mut report = OrderStatusReport::new(
134 account_id,
135 instrument_id,
136 None,
137 venue_order_id,
138 order_side.into(),
139 OrderType::Market,
140 TimeInForce::Gtc,
141 OrderStatus::Filled,
142 order_qty,
143 order_qty,
144 ts_now,
145 ts_now,
146 ts_now,
147 None,
148 )
149 .with_avg_px(avg_px);
150
151 if let Some(venue_position_id) = venue_position_id {
152 report = report.with_venue_position_id(venue_position_id);
153 }
154
155 Some(report)
156}
157
158#[derive(Debug, Clone, PartialEq, Eq)]
160pub(crate) enum ReportClientCoverage {
161 Resolved(IndexSet<ClientId>),
162 Unavailable(IndexSet<ClientId>),
163 Unresolved,
164}
165
166#[derive(Debug, Clone)]
168pub struct ExternalOrderMetadata {
169 pub client_order_id: ClientOrderId,
170 pub venue_order_id: VenueOrderId,
171 pub instrument_id: InstrumentId,
172 pub strategy_id: StrategyId,
173 pub ts_init: UnixNanos,
174}
175
176#[derive(Debug, Default)]
178pub struct ReconciliationResult {
179 pub events: Vec<OrderEventAny>,
181 pub external_orders: Vec<ExternalOrderMetadata>,
183}
184
185#[derive(Debug, Default)]
187pub struct InflightCheckResult {
188 pub events: Vec<OrderEventAny>,
190 pub queries: Vec<TradingCommand>,
192}
193
194#[derive(Debug, Default)]
195pub(crate) struct OpenOrderReconciliationResult {
196 pub events: Vec<OrderEventAny>,
197 pub targeted_queries: Vec<TargetedOrderQuery>,
198}
199
200#[derive(Debug, Clone)]
201pub(crate) struct TargetedOrderQuery {
202 client_order_id: ClientOrderId,
203 responsible_clients: IndexSet<ClientId>,
204 command: GenerateOrderStatusReport,
205}
206
207impl TargetedOrderQuery {
208 #[cfg(feature = "node")]
209 pub(crate) const fn client_order_id(&self) -> ClientOrderId {
210 self.client_order_id
211 }
212}
213
214#[derive(Debug)]
215pub(crate) struct TargetedOrderReportResult {
216 client_order_id: ClientOrderId,
217 client_id: Option<ClientId>,
218 report: Option<OrderStatusReport>,
219 coverage_complete: bool,
220}
221
222#[derive(Debug)]
223pub(crate) struct SourcedOrderStatusReport {
224 pub client_id: ClientId,
225 pub report: OrderStatusReport,
226}
227
228#[derive(Debug, Clone)]
230pub(crate) struct OpenOrderReportCheck {
231 pub command: GenerateOrderStatusReports,
232 pub filtered_orders: Vec<OrderAny>,
233 pub client_coverage: IndexMap<ClientOrderId, ReportClientCoverage>,
234 pub start: Option<UnixNanos>,
235}
236
237#[derive(Debug, Clone)]
239pub(crate) struct PositionReportCheck {
240 pub command: GeneratePositionStatusReports,
241 pub client_coverage: IndexMap<InstrumentAccountKey, ReportClientCoverage>,
242 pub activity_revisions: IndexMap<InstrumentAccountKey, u64>,
243}
244
245#[cfg(feature = "node")]
246#[derive(Debug)]
247pub(crate) struct PositionFillReportQuery {
248 pub key: InstrumentAccountKey,
249 pub client_id: ClientId,
250 pub command: GenerateFillReports,
251}
252
253#[cfg(feature = "node")]
254#[derive(Debug)]
255pub(crate) struct PositionFillReportPlan {
256 pub queries: Vec<PositionFillReportQuery>,
257 pub discrepancy_keys: IndexSet<InstrumentAccountKey>,
258}
259
260#[cfg(feature = "node")]
261#[derive(Debug)]
262pub(crate) enum PositionFillReportPreparation {
263 Ready,
264 InferredOverlap,
265 Unattributed,
266}
267
268struct PositionQuantityComparison {
269 cached_positions: Vec<Position>,
270 cached_signed_qty: Decimal,
271 cached_long_qty: Decimal,
272 cached_short_qty: Decimal,
273 venue_signed_qty: Decimal,
274 venue_long_qty: Decimal,
275 venue_short_qty: Decimal,
276 nonflat_count: usize,
277 venue_report: Option<PositionStatusReport>,
278 venue_has_side_reports: bool,
279}
280
281impl PositionQuantityComparison {
282 fn quantities_match(&self, tolerance: Decimal) -> bool {
283 let net_qty_matches = (self.cached_signed_qty - self.venue_signed_qty).abs() <= tolerance;
284 let side_qty_matches = (self.cached_long_qty - self.venue_long_qty).abs() <= tolerance
285 && (self.cached_short_qty - self.venue_short_qty).abs() <= tolerance;
286
287 net_qty_matches && (!self.venue_has_side_reports || side_qty_matches)
288 }
289
290 fn report_shape(&self) -> PositionReportShape {
291 if self.nonflat_count > 1 || self.venue_has_side_reports {
292 PositionReportShape::MultiLeg
293 } else {
294 PositionReportShape::Unambiguous
295 }
296 }
297}
298
299struct RetainedFillState {
300 fill_keys: IndexSet<(AccountId, InstrumentId, TradeId)>,
301 missing_order_ids: IndexSet<(AccountId, InstrumentId, ClientOrderId)>,
302 missing_venue_order_ids: IndexSet<(AccountId, InstrumentId, VenueOrderId)>,
303 netting_lifecycle_starts: IndexMap<AccountInstrumentStrategyKey, UnixNanos>,
304}
305
306struct HistoricalFillGroup {
307 venue_order_id: VenueOrderId,
308 account_id: AccountId,
309 instrument_id: InstrumentId,
310 strategy_id: StrategyId,
311 order_side: OrderSide,
312 quantity: Decimal,
313 reduce_only: bool,
314 ts_event: UnixNanos,
315 ts_last: UnixNanos,
316}
317
318#[derive(Default)]
319struct ReconciliationFillQueue {
320 pending_fill_keys: IndexSet<FillKey>,
321 event_fill_keys: IndexMap<UUID4, FillKey>,
322}
323
324impl ReconciliationFillQueue {
325 fn push(&mut self, events: &mut Vec<OrderEventAny>, event: OrderEventAny, fill_key: FillKey) {
326 let OrderEventAny::Filled(fill) = &event else {
327 unreachable!("reported fills always create filled events");
328 };
329
330 self.pending_fill_keys.insert(fill_key);
331 self.event_fill_keys.insert(fill.event_id, fill_key);
332 events.push(event);
333 }
334}
335
336#[expect(
338 clippy::struct_excessive_bools,
339 reason = "config flags mirror the live execution engine configuration surface"
340)]
341#[derive(Debug, Clone)]
342pub struct ExecutionManagerConfig {
343 pub trader_id: TraderId,
345 pub reconciliation: bool,
347 pub lookback_mins: Option<u64>,
349 pub reconciliation_instrument_ids: IndexSet<InstrumentId>,
351 pub filter_unclaimed_external: bool,
353 pub filter_position_reports: bool,
355 pub filtered_client_order_ids: IndexSet<ClientOrderId>,
357 pub generate_missing_orders: bool,
359 pub inflight_check_interval_ms: u32,
361 pub inflight_threshold_ms: u64,
363 pub inflight_max_retries: u32,
365 pub open_check_interval_secs: Option<f64>,
367 pub open_check_lookback_mins: Option<u64>,
369 pub open_check_threshold_ns: u64,
371 pub open_check_missing_retries: u32,
373 pub open_check_open_only: bool,
375 pub max_single_order_queries_per_cycle: u32,
377 pub single_order_query_delay_ms: u32,
379 pub position_check_interval_secs: Option<f64>,
381 pub position_check_lookback_mins: u64,
383 pub position_check_threshold_ns: u64,
385 pub position_check_retries: u32,
387 pub purge_closed_orders_buffer_mins: Option<u32>,
389 pub purge_closed_positions_buffer_mins: Option<u32>,
391 pub purge_account_events_lookback_mins: Option<u32>,
393 pub purge_from_database: bool,
395}
396
397impl Default for ExecutionManagerConfig {
398 fn default() -> Self {
399 Self {
400 trader_id: TraderId::default(),
401 reconciliation: true,
402 lookback_mins: Some(60),
403 reconciliation_instrument_ids: IndexSet::new(),
404 filter_unclaimed_external: false,
405 filter_position_reports: false,
406 filtered_client_order_ids: IndexSet::new(),
407 generate_missing_orders: true,
408 inflight_check_interval_ms: 2_000,
409 inflight_threshold_ms: 5_000,
410 inflight_max_retries: 5,
411 open_check_interval_secs: None,
412 open_check_lookback_mins: Some(60),
413 open_check_threshold_ns: 5_000_000_000,
414 open_check_missing_retries: 5,
415 open_check_open_only: true,
416 max_single_order_queries_per_cycle: 5,
417 single_order_query_delay_ms: 100,
418 position_check_interval_secs: None,
419 position_check_lookback_mins: 60,
420 position_check_threshold_ns: 60_000_000_000,
421 position_check_retries: 3,
422 purge_closed_orders_buffer_mins: None,
423 purge_closed_positions_buffer_mins: None,
424 purge_account_events_lookback_mins: None,
425 purge_from_database: false,
426 }
427 }
428}
429
430impl ExecutionManagerConfig {
431 pub fn validate(&self) -> ConfigResult<()> {
438 let mut errors = ConfigErrorCollector::with_capacity(3);
439
440 if let Some(mins) = self.lookback_mins {
441 errors.check(
442 checked_mins_to_secs(mins).is_some(),
443 ConfigError::range(
444 "ExecutionManagerConfig.lookback_mins",
445 format!("{mins} minutes (must fit in `u64` seconds)"),
446 ),
447 );
448 }
449
450 if let Some(mins) = self.open_check_lookback_mins {
451 errors.check(
452 checked_mins_to_nanos(mins).is_some(),
453 ConfigError::range(
454 "ExecutionManagerConfig.open_check_lookback_mins",
455 format!("{mins} minutes (must fit in `u64` nanoseconds)"),
456 ),
457 );
458 }
459
460 errors.check(
461 checked_mins_to_nanos(self.position_check_lookback_mins).is_some(),
462 ConfigError::range(
463 "ExecutionManagerConfig.position_check_lookback_mins",
464 format!(
465 "{} minutes (must fit in `u64` nanoseconds)",
466 self.position_check_lookback_mins
467 ),
468 ),
469 );
470
471 errors.into_result()
472 }
473
474 #[must_use]
476 pub fn with_trader_id(mut self, trader_id: TraderId) -> Self {
477 self.trader_id = trader_id;
478 self
479 }
480}
481
482#[derive(Debug, Clone)]
484struct InflightCheck {
485 #[allow(dead_code)]
486 pub client_order_id: ClientOrderId,
487 pub submitted_at: dst::time::Instant,
488 pub retry_count: u32,
489 pub last_query_at: Option<dst::time::Instant>,
492}
493
494#[derive(Clone, Copy, PartialEq, Eq)]
495enum PositionReportShape {
496 Unambiguous,
497 MultiLeg,
498}
499
500#[derive(Clone, Copy)]
501struct PositionReconciliationState {
502 report_shape: PositionReportShape,
503 retries: u32,
504}
505
506#[derive(Clone)]
529pub struct ExecutionManager {
530 clock: Rc<RefCell<dyn Clock>>,
531 cache: Rc<RefCell<Cache>>,
532 config: ExecutionManagerConfig,
533 inflight_checks: IndexMap<ClientOrderId, InflightCheck>,
534 external_order_claims: IndexMap<InstrumentId, StrategyId>,
535 processed_fills: RecencyMap<FillKey>,
536 recon_check_retries: IndexMap<ClientOrderId, u32>,
537 order_query_recency: RecencyMap<ClientOrderId>,
538 order_local_activity: RecencyMap<ClientOrderId>,
539 position_local_activity: RecencyMap<InstrumentAccountKey>,
541 position_local_activity_revisions: IndexMap<InstrumentAccountKey, u64>,
542 position_reconciliation_states: IndexMap<InstrumentAccountKey, PositionReconciliationState>,
543 position_reconciliation_tolerances: IndexMap<AccountId, Decimal>,
544 recent_fills_cache: RecencyMap<FillKey>,
545 missing_order_coverage_warnings: IndexSet<ClientOrderId>,
546 open_check_lookback_warnings: IndexSet<ClientOrderId>,
547 unresolved_order_coverage: IndexSet<ClientOrderId>,
548 targeted_order_queries: IndexSet<ClientOrderId>,
549}
550
551impl Debug for ExecutionManager {
552 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
553 f.debug_struct(stringify!(ExecutionManager))
554 .field("config", &self.config)
555 .field("inflight_checks", &self.inflight_checks)
556 .field("external_order_claims", &self.external_order_claims)
557 .field("processed_fills", &self.processed_fills)
558 .field("recon_check_retries", &self.recon_check_retries)
559 .finish_non_exhaustive()
560 }
561}
562
563impl ExecutionManager {
564 pub fn new(
570 clock: Rc<RefCell<dyn Clock>>,
571 cache: Rc<RefCell<Cache>>,
572 config: ExecutionManagerConfig,
573 ) -> ConfigResult<Self> {
574 config.validate()?;
575
576 Ok(Self {
577 clock,
578 cache,
579 config,
580 inflight_checks: IndexMap::new(),
581 external_order_claims: IndexMap::new(),
582 processed_fills: RecencyMap::default(),
583 recon_check_retries: IndexMap::new(),
584 order_query_recency: RecencyMap::default(),
585 order_local_activity: RecencyMap::default(),
586 position_local_activity: RecencyMap::default(),
587 position_local_activity_revisions: IndexMap::new(),
588 position_reconciliation_states: IndexMap::new(),
589 position_reconciliation_tolerances: IndexMap::new(),
590 recent_fills_cache: RecencyMap::default(),
591 missing_order_coverage_warnings: IndexSet::new(),
592 open_check_lookback_warnings: IndexSet::new(),
593 unresolved_order_coverage: IndexSet::new(),
594 targeted_order_queries: IndexSet::new(),
595 })
596 }
597
598 pub(crate) fn set_position_reconciliation_tolerance(
599 &mut self,
600 account_id: AccountId,
601 tolerance: Decimal,
602 ) {
603 let tolerance = if tolerance < Decimal::ZERO {
604 log::error!(
605 "Invalid negative position reconciliation tolerance {tolerance} for \
606 {account_id}; using the default"
607 );
608 DEFAULT_POSITION_RECONCILIATION_TOLERANCE
609 } else {
610 tolerance
611 };
612 self.position_reconciliation_tolerances
613 .insert(account_id, tolerance);
614 }
615
616 fn position_reconciliation_tolerance(&self, account_id: AccountId) -> Decimal {
617 self.position_reconciliation_tolerances
618 .get(&account_id)
619 .copied()
620 .unwrap_or(DEFAULT_POSITION_RECONCILIATION_TOLERANCE)
621 }
622
623 #[allow(unknown_lints, reason = "Clippy lint is unavailable on Rust 1.97")]
629 #[expect(
630 clippy::unused_async,
631 clippy::unused_async_trait_impl,
632 reason = "public reconciliation API stays async; live node and test callers await it"
633 )]
634 pub async fn reconcile_execution_mass_status(
635 &mut self,
636 mass_status: ExecutionMassStatus,
637 exec_engine: Rc<RefCell<ExecutionEngine>>,
638 ) -> ReconciliationResult {
639 if exec_engine
640 .borrow()
641 .get_client(&mass_status.client_id)
642 .is_none()
643 {
644 log::error!(
645 "Cannot reconcile ExecutionMassStatus from unknown client {}",
646 mass_status.client_id
647 );
648 return ReconciliationResult::default();
649 }
650
651 self.validate_mass_status_order_sources(&mass_status);
652
653 let raw_order_status_topic =
658 MessagingSwitchboard::reconciliation_raw_order_status_report_topic();
659
660 for report in mass_status.order_reports().values() {
661 msgbus::publish_any(raw_order_status_topic, report);
662 }
663
664 let raw_fill_topic = MessagingSwitchboard::reconciliation_raw_fill_report_topic();
665
666 for fills in mass_status.fill_reports().values() {
667 for fill in fills {
668 msgbus::publish_any(raw_fill_topic, fill);
669 }
670 }
671
672 let raw_position_topic =
673 MessagingSwitchboard::reconciliation_raw_position_status_report_topic();
674
675 for reports in mass_status.position_reports().values() {
676 for report in reports {
677 msgbus::publish_any(raw_position_topic, report);
678 }
679 }
680
681 if exec_engine
682 .borrow()
683 .get_client(&mass_status.client_id)
684 .is_none()
685 {
686 log::error!(
687 "Execution client {} disappeared while publishing raw mass status reports",
688 mass_status.client_id
689 );
690 return ReconciliationResult::default();
691 }
692
693 let venue = mass_status.venue;
694 let order_count = mass_status.order_reports().len();
695 let fill_count: usize = mass_status.fill_reports().values().map(Vec::len).sum();
696 let position_count: usize = mass_status.position_reports().values().map(Vec::len).sum();
697
698 log_info!(
699 "Reconciling ExecutionMassStatus for {venue}",
700 color = LogColor::Blue
701 );
702 log_info!(
703 "Received {order_count} order(s), {fill_count} fill(s), {position_count} position(s)",
704 color = LogColor::Blue
705 );
706
707 let retained_fill_state = self.retained_fill_state();
708 let reported_fill_keys: IndexSet<(AccountId, InstrumentId, TradeId)> = mass_status
709 .fill_reports()
710 .values()
711 .flatten()
712 .map(|fill| (fill.account_id, fill.instrument_id, fill.trade_id))
713 .collect();
714 let (adjusted_order_reports, adjusted_fill_reports) =
715 self.adjust_mass_status_fills(&mass_status);
716 let order_only_venue_order_ids = self.order_only_venue_order_ids(
717 &mass_status,
718 &adjusted_order_reports,
719 &adjusted_fill_reports,
720 &retained_fill_state,
721 );
722
723 let mut events = Vec::new();
724 let mut external_orders = Vec::new();
725 let mut orders_reconciled = 0usize;
726 let mut external_orders_created = 0usize;
727 let mut open_orders_initialized = 0usize;
728 let mut orders_skipped_no_instrument = 0usize;
729 let mut orders_skipped_duplicate = 0usize;
730 let mut fills_applied = 0usize;
731 let mut fill_queue = ReconciliationFillQueue::default();
732
733 let fill_reports = &adjusted_fill_reports;
734 let mut seen_fill_keys: IndexSet<FillKey> = IndexSet::new();
735
736 for fills in fill_reports.values() {
737 for fill in fills {
738 let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
739 if !seen_fill_keys.insert(fill_key) {
740 log::warn!(
741 "Duplicate trade_id {} for {} in mass status",
742 fill.trade_id,
743 fill.instrument_id
744 );
745 }
746 }
747 }
748
749 let order_reports = Self::deduplicate_order_reports(adjusted_order_reports.values());
751 let mut orders_skipped_filtered = 0usize;
752
753 for report in order_reports.values() {
754 if self.should_skip_order_report(report) {
755 orders_skipped_filtered += 1;
756 continue;
757 }
758
759 if let Some(client_order_id) = &report.client_order_id {
760 if let Some(cached_order) = self.get_order(*client_order_id)
761 && Self::is_exact_order_match(&cached_order, report)
762 {
763 log::debug!("Skipping order {client_order_id}: already in sync with venue");
764 orders_skipped_duplicate += 1;
765
766 if let Err(e) = self
768 .cache
769 .borrow_mut()
770 .index_venue_order_id(client_order_id, &report.venue_order_id)
771 {
772 log::warn!("Failed to index venue order ID: {e}");
773 }
774
775 continue;
776 }
777
778 if let Some(cached_order) = self.get_order(*client_order_id)
780 && cached_order.is_closed()
781 && cached_order
782 .tags()
783 .is_some_and(|tags| tags.contains(&*TAG_RECONCILIATION))
784 {
785 log::debug!(
786 "Skipping closed reconciliation order {client_order_id}: \
787 synthetic position adjustment from previous session",
788 );
789 orders_skipped_duplicate += 1;
790 continue;
791 }
792
793 if let Some(order) = self.get_order(*client_order_id) {
794 let instrument = self.get_instrument(&report.instrument_id);
795 log::info!(
796 color = LogColor::Blue as u8;
797 "Reconciling {} {} {} [{}] -> [{}]",
798 client_order_id,
799 report.venue_order_id,
800 report.instrument_id,
801 order.status(),
802 report.order_status,
803 );
804
805 let order_fills: Vec<&FillReport> = fill_reports
806 .get(&report.venue_order_id)
807 .map(|f| f.iter().collect())
808 .unwrap_or_default();
809 let engine_ref = exec_engine.borrow();
810 let commission_client = engine_ref.get_client(&mass_status.client_id);
811 let order_events = self.reconcile_order_with_fills(
812 &order,
813 report,
814 &order_fills,
815 instrument.as_ref(),
816 &mut fill_queue,
817 commission_client,
818 );
819 drop(engine_ref);
820
821 if !order_events.is_empty() {
822 orders_reconciled += 1;
823 fills_applied += order_events
824 .iter()
825 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
826 .count();
827 events.extend(order_events);
828 }
829
830 if let Err(e) = self
832 .cache
833 .borrow_mut()
834 .index_venue_order_id(client_order_id, &report.venue_order_id)
835 {
836 log::warn!("Failed to index venue order ID: {e}");
837 }
838 } else if let Some(order) = self.get_order_by_venue_order_id(report.venue_order_id)
839 {
840 let instrument = self.get_instrument(&report.instrument_id);
842
843 log::info!(
844 color = LogColor::Blue as u8;
845 "Reconciling {} (matched by venue_order_id {}) {} [{}] -> [{}]",
846 order.client_order_id(),
847 report.venue_order_id,
848 report.instrument_id,
849 order.status(),
850 report.order_status,
851 );
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 order_events = self.reconcile_order_with_fills(
860 &order,
861 report,
862 &order_fills,
863 instrument.as_ref(),
864 &mut fill_queue,
865 commission_client,
866 );
867 drop(engine_ref);
868
869 if !order_events.is_empty() {
870 orders_reconciled += 1;
871 fills_applied += order_events
872 .iter()
873 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
874 .count();
875 events.extend(order_events);
876 }
877
878 if let Err(e) = self
879 .cache
880 .borrow_mut()
881 .index_venue_order_id(&order.client_order_id(), &report.venue_order_id)
882 {
883 log::warn!("Failed to index venue order ID: {e}");
884 }
885 } else if !self.config.filter_unclaimed_external {
886 if let Some(instrument) = self.get_instrument(&report.instrument_id) {
887 let order_fills: Vec<&FillReport> = fill_reports
888 .get(&report.venue_order_id)
889 .map(|f| f.iter().collect())
890 .unwrap_or_default();
891 let engine_ref = exec_engine.borrow();
892 let commission_client = engine_ref.get_client(&mass_status.client_id);
893 let (external_events, metadata) = self.handle_external_order(
894 report,
895 mass_status.account_id,
896 &instrument,
897 &order_fills,
898 false, Some(&mut fill_queue),
900 commission_client,
901 );
902 drop(engine_ref);
903
904 if !external_events.is_empty() {
905 external_orders_created += 1;
906 fills_applied += external_events
907 .iter()
908 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
909 .count();
910
911 if report.order_status.is_open() {
912 open_orders_initialized += 1;
913 }
914
915 events.extend(external_events);
916
917 if let Some(m) = metadata {
918 external_orders.push(m);
919 }
920 }
921 } else {
922 orders_skipped_no_instrument += 1;
923 }
924 }
925 } else if let Some(order) = self.get_order_by_venue_order_id(report.venue_order_id) {
926 let instrument = self.get_instrument(&report.instrument_id);
928 log::info!(
929 color = LogColor::Blue as u8;
930 "Reconciling {} (matched by venue_order_id {}) {} [{}] -> [{}]",
931 order.client_order_id(),
932 report.venue_order_id,
933 report.instrument_id,
934 order.status(),
935 report.order_status,
936 );
937
938 let order_fills: Vec<&FillReport> = fill_reports
939 .get(&report.venue_order_id)
940 .map(|f| f.iter().collect())
941 .unwrap_or_default();
942 let engine_ref = exec_engine.borrow();
943 let commission_client = engine_ref.get_client(&mass_status.client_id);
944 let order_events = self.reconcile_order_with_fills(
945 &order,
946 report,
947 &order_fills,
948 instrument.as_ref(),
949 &mut fill_queue,
950 commission_client,
951 );
952 drop(engine_ref);
953
954 if !order_events.is_empty() {
955 orders_reconciled += 1;
956 fills_applied += order_events
957 .iter()
958 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
959 .count();
960 events.extend(order_events);
961 }
962
963 if let Err(e) = self
964 .cache
965 .borrow_mut()
966 .index_venue_order_id(&order.client_order_id(), &report.venue_order_id)
967 {
968 log::warn!("Failed to index venue order ID: {e}");
969 }
970 } else if let Some(instrument) = self.get_instrument(&report.instrument_id) {
971 let is_synthetic = report.venue_order_id.as_str().starts_with("S-");
973
974 let order_fills: Vec<&FillReport> = fill_reports
975 .get(&report.venue_order_id)
976 .map(|f| f.iter().collect())
977 .unwrap_or_default();
978 let engine_ref = exec_engine.borrow();
979 let commission_client = engine_ref.get_client(&mass_status.client_id);
980 let (external_events, metadata) = self.handle_external_order(
981 report,
982 mass_status.account_id,
983 &instrument,
984 &order_fills,
985 is_synthetic,
986 Some(&mut fill_queue),
987 commission_client,
988 );
989 drop(engine_ref);
990
991 if !external_events.is_empty() {
992 external_orders_created += 1;
993 fills_applied += external_events
994 .iter()
995 .filter(|e| matches!(e, OrderEventAny::Filled(_)))
996 .count();
997
998 if report.order_status.is_open() {
999 open_orders_initialized += 1;
1000 }
1001
1002 events.extend(external_events);
1003
1004 if let Some(m) = metadata {
1005 external_orders.push(m);
1006 }
1007 }
1008 } else {
1009 orders_skipped_no_instrument += 1;
1010 }
1011 }
1012
1013 let processed_venue_order_ids: IndexSet<VenueOrderId> =
1015 order_reports.keys().copied().collect();
1016
1017 for (venue_order_id, fills) in fill_reports {
1018 if processed_venue_order_ids.contains(venue_order_id) {
1019 continue;
1020 }
1021
1022 let Some(first_fill) = fills.first() else {
1023 continue;
1024 };
1025
1026 if !self.should_reconcile_instrument(&first_fill.instrument_id) {
1027 log::debug!(
1028 "Skipping orphan fills for {}: not in reconciliation_instrument_ids",
1029 first_fill.instrument_id
1030 );
1031 continue;
1032 }
1033
1034 if let Some(client_order_id) = &first_fill.client_order_id
1036 && self
1037 .config
1038 .filtered_client_order_ids
1039 .contains(client_order_id)
1040 {
1041 log::debug!(
1042 "Skipping orphan fills for {client_order_id}: in filtered_client_order_ids"
1043 );
1044 continue;
1045 }
1046
1047 let order = first_fill
1048 .client_order_id
1049 .as_ref()
1050 .and_then(|id| self.get_order(*id))
1051 .or_else(|| self.get_order_by_venue_order_id(*venue_order_id));
1052
1053 if let Some(ref order) = order
1055 && self
1056 .config
1057 .filtered_client_order_ids
1058 .contains(&order.client_order_id())
1059 {
1060 log::debug!(
1061 "Skipping orphan fills for {}: in filtered_client_order_ids",
1062 order.client_order_id()
1063 );
1064 continue;
1065 }
1066
1067 if let Some(order) = order {
1068 let instrument_id = order.instrument_id();
1069 if let Some(instrument) = self.get_instrument(&instrument_id) {
1070 let mut sorted_fills: Vec<&FillReport> = fills.iter().collect();
1071 sorted_fills.sort_by_key(|f| f.ts_event);
1072
1073 for fill in sorted_fills {
1074 if let Some((event, fill_key)) = self.create_order_fill(
1075 &order,
1076 fill,
1077 &instrument,
1078 &fill_queue.pending_fill_keys,
1079 ) {
1080 fills_applied += 1;
1081 fill_queue.push(&mut events, event, fill_key);
1082 }
1083 }
1084 } else {
1085 orders_skipped_no_instrument += 1;
1086 }
1087 } else if fills.iter().any(FillReport::has_venue_position_id) {
1088 if !self.config.generate_missing_orders {
1089 log::debug!(
1090 "Skipping orphan fills for venue order {venue_order_id}: \
1091 `generate_missing_orders` is disabled"
1092 );
1093 orders_skipped_filtered += 1;
1094 continue;
1095 }
1096
1097 let Some(instrument) = self.get_instrument(&first_fill.instrument_id) else {
1098 orders_skipped_no_instrument += 1;
1099 continue;
1100 };
1101
1102 let mut sorted_fills: Vec<&FillReport> = fills.iter().collect();
1103 sorted_fills.sort_by_key(|fill| fill.ts_event);
1104
1105 let report = match Self::create_orphan_fill_order_report(&sorted_fills, &instrument)
1106 {
1107 Ok(report) => report,
1108 Err(e) => {
1109 log::error!(
1110 "Cannot materialize orphan fills for venue order {venue_order_id}: {e}"
1111 );
1112
1113 continue;
1114 }
1115 };
1116
1117 let engine_ref = exec_engine.borrow();
1118 let commission_client = engine_ref.get_client(&mass_status.client_id);
1119 let (external_events, metadata) = self.handle_external_order(
1120 &report,
1121 mass_status.account_id,
1122 &instrument,
1123 &sorted_fills,
1124 false,
1125 Some(&mut fill_queue),
1126 commission_client,
1127 );
1128 drop(engine_ref);
1129
1130 if !external_events.is_empty() {
1131 external_orders_created += 1;
1132 fills_applied += external_events
1133 .iter()
1134 .filter(|event| matches!(event, OrderEventAny::Filled(_)))
1135 .count();
1136
1137 events.extend(external_events);
1138
1139 if let Some(metadata) = metadata {
1140 external_orders.push(metadata);
1141 }
1142 }
1143 }
1144 }
1145
1146 events.sort_by_key(OrderEventAny::ts_event);
1147
1148 let mut unapplied_fill_position_ids = IndexSet::new();
1149
1150 for event in &events {
1151 if let OrderEventAny::Filled(fill) = event
1152 && Self::should_project_reconciliation_fill(
1153 fill,
1154 &retained_fill_state,
1155 &reported_fill_keys,
1156 &order_only_venue_order_ids,
1157 )
1158 {
1159 exec_engine.borrow_mut().project_reconciliation_fill(fill);
1160 } else {
1161 exec_engine.borrow_mut().process(event);
1162 }
1163
1164 if let OrderEventAny::Filled(fill) = event
1165 && let Some(fill_key) = fill_queue.event_fill_keys.get(&fill.event_id).copied()
1166 {
1167 if self.is_fill_applied(fill, fill_key) {
1168 self.processed_fills.mark(fill_key);
1169 } else if let Some(venue_position_id) = fill.position_id {
1170 log::error!(
1171 "Skipping reconciliation for venue position {venue_position_id}: historical fill {} was not applied",
1172 fill.trade_id,
1173 );
1174
1175 unapplied_fill_position_ids.insert(venue_position_id);
1176 }
1177 }
1178 }
1179
1180 let mut positions_created = 0usize;
1181
1182 if !self.config.filter_position_reports {
1183 let instruments_with_unattributed_fills: IndexSet<InstrumentId> = mass_status
1186 .fill_reports()
1187 .values()
1188 .flatten()
1189 .filter(|f| f.venue_position_id.is_none())
1190 .map(|f| f.instrument_id)
1191 .chain(
1192 mass_status
1193 .order_reports()
1194 .values()
1195 .filter(|r| !r.filled_qty.is_zero() && r.venue_position_id.is_none())
1196 .map(|r| r.instrument_id),
1197 )
1198 .collect();
1199
1200 for (instrument_id, reports) in mass_status.position_reports() {
1201 if !self.should_reconcile_instrument(&instrument_id) {
1202 log::debug!(
1203 "Skipping position reports for {instrument_id}: not in reconciliation_instrument_ids"
1204 );
1205 continue;
1206 }
1207
1208 for report in reports {
1209 if report.venue_position_id.is_some_and(|venue_position_id| {
1210 unapplied_fill_position_ids.contains(&venue_position_id)
1211 }) {
1212 continue;
1213 }
1214
1215 if let Some(position_events) = self.reconcile_position_report(
1216 &report,
1217 mass_status.account_id,
1218 &instruments_with_unattributed_fills,
1219 ) {
1220 for event in position_events {
1221 exec_engine.borrow_mut().process(&event);
1222 events.push(event);
1223 }
1224
1225 positions_created += 1;
1226 }
1227 }
1228 }
1229 }
1230
1231 if orders_skipped_no_instrument > 0 {
1232 log::warn!("{orders_skipped_no_instrument} orders skipped (instrument not in cache)");
1233 }
1234
1235 if orders_skipped_duplicate > 0 {
1236 log::debug!("{orders_skipped_duplicate} orders skipped (already in sync)");
1237 }
1238
1239 if orders_skipped_filtered > 0 {
1240 log::debug!("{orders_skipped_filtered} orders skipped (filtered by config)");
1241 }
1242
1243 log::info!(
1244 color = LogColor::Blue as u8;
1245 "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}",
1246 );
1247
1248 ReconciliationResult {
1249 events,
1250 external_orders,
1251 }
1252 }
1253
1254 fn create_orphan_fill_order_report(
1255 fills: &[&FillReport],
1256 instrument: &InstrumentAny,
1257 ) -> anyhow::Result<OrderStatusReport> {
1258 let Some(first) = fills.first() else {
1259 anyhow::bail!("fill group is empty");
1260 };
1261 let venue_position_id = first
1262 .venue_position_id
1263 .ok_or_else(|| anyhow::anyhow!("venue position ID is missing"))?;
1264
1265 for fill in fills.iter().skip(1) {
1266 anyhow::ensure!(
1267 fill.account_id == first.account_id,
1268 "account ID differs across fill group"
1269 );
1270 anyhow::ensure!(
1271 fill.instrument_id == first.instrument_id,
1272 "instrument ID differs across fill group"
1273 );
1274 anyhow::ensure!(
1275 fill.venue_order_id == first.venue_order_id,
1276 "venue order ID differs across fill group"
1277 );
1278 anyhow::ensure!(
1279 fill.client_order_id == first.client_order_id,
1280 "client order ID differs across fill group"
1281 );
1282 anyhow::ensure!(
1283 fill.order_side == first.order_side,
1284 "order side differs across fill group"
1285 );
1286 anyhow::ensure!(
1287 fill.venue_position_id == first.venue_position_id,
1288 "venue position ID differs across fill group"
1289 );
1290 }
1291
1292 anyhow::ensure!(
1293 first.instrument_id == instrument.id(),
1294 "instrument metadata does not match fill group"
1295 );
1296
1297 let (quantity, notional) = fills.iter().try_fold(
1298 (Decimal::ZERO, Decimal::ZERO),
1299 |(quantity, notional), fill| {
1300 let fill_quantity = fill.last_qty.as_decimal();
1301 let quantity = quantity.checked_add(fill_quantity).ok_or_else(|| {
1302 anyhow::anyhow!("fill quantity overflow while aggregating fill group")
1303 })?;
1304
1305 let fill_notional = fill_quantity
1306 .checked_mul(fill.last_px.as_decimal())
1307 .ok_or_else(|| {
1308 anyhow::anyhow!("fill notional overflow while aggregating fill group")
1309 })?;
1310
1311 let notional = notional.checked_add(fill_notional).ok_or_else(|| {
1312 anyhow::anyhow!("fill notional overflow while aggregating fill group")
1313 })?;
1314
1315 Ok::<_, anyhow::Error>((quantity, notional))
1316 },
1317 )?;
1318
1319 anyhow::ensure!(
1320 quantity > Decimal::ZERO,
1321 "fill group quantity is not positive"
1322 );
1323
1324 let order_qty = Quantity::from_decimal_dp(quantity, instrument.size_precision())?;
1325 let avg_px = notional
1326 .checked_div(quantity)
1327 .ok_or_else(|| anyhow::anyhow!("fill group average price is not representable"))?;
1328
1329 let ts_accepted = fills
1330 .iter()
1331 .map(|fill| fill.ts_event)
1332 .min()
1333 .expect("non-empty fill group");
1334
1335 let ts_last = fills
1336 .iter()
1337 .map(|fill| fill.ts_event)
1338 .max()
1339 .expect("non-empty fill group");
1340
1341 let ts_init = fills
1342 .iter()
1343 .map(|fill| fill.ts_init)
1344 .max()
1345 .expect("non-empty fill group");
1346
1347 let report = OrderStatusReport::new(
1348 first.account_id,
1349 first.instrument_id,
1350 first.client_order_id,
1351 first.venue_order_id,
1352 first.order_side.into(),
1353 OrderType::Market,
1354 TimeInForce::Gtc,
1355 OrderStatus::Filled,
1356 order_qty,
1357 order_qty,
1358 ts_accepted,
1359 ts_last,
1360 ts_init,
1361 None,
1362 )
1363 .with_avg_px(avg_px)
1364 .with_venue_position_id(venue_position_id);
1365
1366 Ok(report)
1367 }
1368
1369 fn should_project_reconciliation_fill(
1370 fill: &OrderFilled,
1371 retained_fill_state: &RetainedFillState,
1372 reported_fill_keys: &IndexSet<FillKey>,
1373 order_only_venue_order_ids: &IndexSet<VenueOrderId>,
1374 ) -> bool {
1375 let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
1376 if retained_fill_state.fill_keys.contains(&fill_key)
1377 || order_only_venue_order_ids.contains(&fill.venue_order_id)
1378 {
1379 return true;
1380 }
1381
1382 let order_missing = retained_fill_state.missing_order_ids.contains(&(
1383 fill.account_id,
1384 fill.instrument_id,
1385 fill.client_order_id,
1386 )) || retained_fill_state.missing_venue_order_ids.contains(&(
1387 fill.account_id,
1388 fill.instrument_id,
1389 fill.venue_order_id,
1390 ));
1391
1392 if order_missing && !reported_fill_keys.contains(&fill_key) {
1393 return true;
1394 }
1395
1396 retained_fill_state
1397 .netting_lifecycle_starts
1398 .get(&(fill.account_id, fill.instrument_id, fill.strategy_id))
1399 .is_some_and(|ts_opened| fill.ts_event < *ts_opened)
1400 }
1401
1402 fn retained_fill_state(&self) -> RetainedFillState {
1403 let cache = self.cache.borrow();
1404 let positions = cache.positions(None, None, None, None, None);
1405 let mut fill_keys = IndexSet::new();
1406 let mut missing_order_ids = IndexSet::new();
1407 let mut missing_venue_order_ids = IndexSet::new();
1408 let mut netting_lifecycle_starts = IndexMap::new();
1409
1410 for position in positions {
1411 for fill in &position.events {
1412 fill_keys.insert((position.account_id, position.instrument_id, fill.trade_id));
1413 if cache.order(&fill.client_order_id).is_none() {
1414 missing_order_ids.insert((
1415 position.account_id,
1416 position.instrument_id,
1417 fill.client_order_id,
1418 ));
1419 missing_venue_order_ids.insert((
1420 position.account_id,
1421 position.instrument_id,
1422 fill.venue_order_id,
1423 ));
1424 }
1425 }
1426
1427 if cache.oms_type(&position.id) == Some(OmsType::Netting) {
1428 netting_lifecycle_starts.insert(
1429 (
1430 position.account_id,
1431 position.instrument_id,
1432 position.strategy_id,
1433 ),
1434 position.ts_opened,
1435 );
1436 }
1437 }
1438
1439 RetainedFillState {
1440 fill_keys,
1441 missing_order_ids,
1442 missing_venue_order_ids,
1443 netting_lifecycle_starts,
1444 }
1445 }
1446
1447 fn order_only_venue_order_ids(
1448 &self,
1449 mass_status: &ExecutionMassStatus,
1450 order_reports: &IndexMap<VenueOrderId, OrderStatusReport>,
1451 fill_reports: &IndexMap<VenueOrderId, Vec<FillReport>>,
1452 retained_fill_state: &RetainedFillState,
1453 ) -> IndexSet<VenueOrderId> {
1454 if mass_status.lookback_start().is_none() {
1455 return IndexSet::new();
1456 }
1457
1458 let expected_quantities: IndexMap<AccountInstrumentKey, Decimal> =
1459 if mass_status.reports_complete() {
1460 mass_status
1461 .position_reports()
1462 .into_iter()
1463 .filter_map(|(instrument_id, reports)| {
1464 let [report] = reports.as_slice() else {
1465 return None;
1466 };
1467 report.venue_position_id.is_none().then_some((
1468 (report.account_id, instrument_id),
1469 report.signed_decimal_qty,
1470 ))
1471 })
1472 .collect()
1473 } else {
1474 IndexMap::new()
1475 };
1476 let candidate_instruments: IndexSet<InstrumentId> = order_reports
1477 .values()
1478 .filter(|report| !report.filled_qty.is_zero())
1479 .map(|report| report.instrument_id)
1480 .chain(
1481 fill_reports
1482 .values()
1483 .flatten()
1484 .map(|fill| fill.instrument_id),
1485 )
1486 .collect();
1487
1488 if candidate_instruments.is_empty() {
1489 return IndexSet::new();
1490 }
1491
1492 let mut venue_order_ids: IndexSet<VenueOrderId> = order_reports
1493 .iter()
1494 .filter(|(_, report)| {
1495 candidate_instruments.contains(&report.instrument_id)
1496 && !report.filled_qty.is_zero()
1497 })
1498 .map(|(venue_order_id, _)| *venue_order_id)
1499 .collect();
1500 venue_order_ids.extend(fill_reports.iter().filter_map(|(venue_order_id, fills)| {
1501 fills
1502 .first()
1503 .is_some_and(|fill| candidate_instruments.contains(&fill.instrument_id))
1504 .then_some(*venue_order_id)
1505 }));
1506
1507 if !mass_status.reports_complete() {
1508 log::error!(
1509 "Bounded reconciliation report set is incomplete; projecting {} historical order(s) without position or portfolio effects",
1510 venue_order_ids.len(),
1511 );
1512
1513 return venue_order_ids;
1514 }
1515
1516 let mut order_only = IndexSet::new();
1517 let mut groups = Vec::new();
1518
1519 for venue_order_id in venue_order_ids {
1520 let report = order_reports.get(&venue_order_id);
1521 let fills = fill_reports.get(&venue_order_id);
1522
1523 if report.and_then(|report| report.venue_position_id).is_some()
1524 || fills.is_some_and(|fills| fills.iter().any(FillReport::has_venue_position_id))
1525 {
1526 continue;
1527 }
1528
1529 let cached_order = report
1530 .and_then(|report| report.client_order_id)
1531 .and_then(|client_order_id| self.get_order(client_order_id))
1532 .or_else(|| self.get_order_by_venue_order_id(venue_order_id));
1533 let account_id = report
1534 .map(|report| report.account_id)
1535 .or_else(|| fills.and_then(|fills| fills.first().map(|fill| fill.account_id)));
1536 let instrument_id = report
1537 .map(|report| report.instrument_id)
1538 .or_else(|| fills.and_then(|fills| fills.first().map(|fill| fill.instrument_id)));
1539 let order_side = report
1540 .and_then(|report| report.order_side)
1541 .or_else(|| fills.and_then(|fills| fills.first().map(|fill| fill.order_side)));
1542 let (Some(account_id), Some(instrument_id), Some(order_side)) =
1543 (account_id, instrument_id, order_side)
1544 else {
1545 order_only.insert(venue_order_id);
1546 continue;
1547 };
1548
1549 let coherent_fills = fills.is_none_or(|fills| {
1550 fills.iter().all(|fill| {
1551 fill.account_id == account_id
1552 && fill.instrument_id == instrument_id
1553 && fill.order_side == order_side
1554 })
1555 });
1556 let coherent_cached_order = cached_order.as_ref().is_none_or(|order| {
1557 order.instrument_id() == instrument_id
1558 && order.order_side() == order_side
1559 && order.account_id().is_none_or(|id| id == account_id)
1560 });
1561
1562 if !coherent_fills
1563 || !coherent_cached_order
1564 || (report.is_none() && cached_order.is_none())
1565 {
1566 order_only.insert(venue_order_id);
1567 continue;
1568 }
1569
1570 let strategy_id = cached_order.as_ref().map_or_else(
1571 || {
1572 self.external_order_claims
1573 .get(&instrument_id)
1574 .copied()
1575 .unwrap_or_else(|| StrategyId::from("EXTERNAL"))
1576 },
1577 Order::strategy_id,
1578 );
1579 let reduce_only = report.is_some_and(|report| report.reduce_only)
1580 || cached_order.as_ref().is_some_and(Order::is_reduce_only);
1581
1582 let cached_filled_qty = cached_order
1583 .as_ref()
1584 .map_or(Decimal::ZERO, |order| order.filled_qty().as_decimal());
1585 let reported_fill_qty = fills.map_or(Decimal::ZERO, |fills| {
1586 fills.iter().map(|fill| fill.last_qty.as_decimal()).sum()
1587 });
1588
1589 let unretained_fills: Vec<&FillReport> = fills
1590 .into_iter()
1591 .flatten()
1592 .filter(|fill| {
1593 !retained_fill_state.fill_keys.contains(&(
1594 fill.account_id,
1595 fill.instrument_id,
1596 fill.trade_id,
1597 ))
1598 })
1599 .collect();
1600 let unretained_fill_qty: Decimal = unretained_fills
1601 .iter()
1602 .map(|fill| fill.last_qty.as_decimal())
1603 .sum();
1604
1605 let inferred_qty = report.map_or(Decimal::ZERO, |report| {
1606 (report.filled_qty.as_decimal() - cached_filled_qty - reported_fill_qty)
1607 .max(Decimal::ZERO)
1608 });
1609 let quantity = unretained_fill_qty + inferred_qty;
1610
1611 if quantity.is_zero() {
1612 continue;
1613 }
1614
1615 let inferred_ts = (!inferred_qty.is_zero())
1616 .then(|| report.map(|report| report.ts_last))
1617 .flatten();
1618 let ts_event = unretained_fills
1619 .iter()
1620 .map(|fill| fill.ts_event)
1621 .chain(inferred_ts)
1622 .min()
1623 .unwrap_or(mass_status.ts_init);
1624 let ts_last = unretained_fills
1625 .iter()
1626 .map(|fill| fill.ts_event)
1627 .chain(inferred_ts)
1628 .max()
1629 .unwrap_or(mass_status.ts_init);
1630
1631 groups.push(HistoricalFillGroup {
1632 venue_order_id,
1633 account_id,
1634 instrument_id,
1635 strategy_id,
1636 order_side,
1637 quantity,
1638 reduce_only,
1639 ts_event,
1640 ts_last,
1641 });
1642 }
1643
1644 groups.sort_by_key(|group| group.ts_event);
1645
1646 let mut quantities: IndexMap<AccountInstrumentStrategyKey, Option<Decimal>> =
1647 IndexMap::new();
1648 let mut group_ids: IndexMap<AccountInstrumentStrategyKey, Vec<VenueOrderId>> =
1649 IndexMap::new();
1650 let mut interval_ends: IndexMap<AccountInstrumentStrategyKey, UnixNanos> = IndexMap::new();
1651 let mut ambiguous_keys = IndexSet::new();
1652
1653 for group in &groups {
1654 let key = (group.account_id, group.instrument_id, group.strategy_id);
1655
1656 if interval_ends
1657 .get(&key)
1658 .is_some_and(|end| group.ts_event <= *end)
1659 {
1660 ambiguous_keys.insert(key);
1661 }
1662
1663 interval_ends
1664 .entry(key)
1665 .and_modify(|end| *end = (*end).max(group.ts_last))
1666 .or_insert(group.ts_last);
1667 }
1668
1669 if !ambiguous_keys.is_empty() {
1670 log::error!(
1671 "Bounded reconciliation contains interleaved order fills for {} position key(s); projecting their historical order state only",
1672 ambiguous_keys.len(),
1673 );
1674 }
1675
1676 for group in groups {
1677 let key = (group.account_id, group.instrument_id, group.strategy_id);
1678 group_ids.entry(key).or_default().push(group.venue_order_id);
1679 if ambiguous_keys.contains(&key) {
1680 order_only.insert(group.venue_order_id);
1681 continue;
1682 }
1683 let current_qty = quantities.entry(key).or_insert_with(|| {
1684 let cache = self.cache.borrow();
1685 let positions = cache.positions_open(
1686 None,
1687 Some(&group.instrument_id),
1688 Some(&group.strategy_id),
1689 Some(&group.account_id),
1690 None,
1691 );
1692
1693 if positions.len() > 1
1694 || positions.first().is_some_and(|position| {
1695 cache.oms_type(&position.id) != Some(OmsType::Netting)
1696 })
1697 {
1698 None
1699 } else {
1700 Some(
1701 positions
1702 .first()
1703 .map_or(Decimal::ZERO, |position| position.signed_decimal_qty()),
1704 )
1705 }
1706 });
1707 let Some(current_qty) = current_qty else {
1708 order_only.insert(group.venue_order_id);
1709 continue;
1710 };
1711 let signed_fill_qty = match group.order_side {
1712 OrderSide::Buy => group.quantity,
1713 OrderSide::Sell => -group.quantity,
1714 };
1715 let reduces = !current_qty.is_zero()
1716 && current_qty.is_sign_negative() != signed_fill_qty.is_sign_negative()
1717 && group.quantity <= current_qty.abs();
1718 if group.reduce_only && !reduces {
1719 log::warn!(
1720 "Cannot apply bounded reduce-only order {} for {} without a coherent predecessor; projecting order state only",
1721 group.venue_order_id,
1722 group.instrument_id,
1723 );
1724 order_only.insert(group.venue_order_id);
1725 continue;
1726 }
1727 *current_qty += signed_fill_qty;
1728 }
1729
1730 let mut keys_by_position: IndexMap<
1731 AccountInstrumentKey,
1732 Vec<AccountInstrumentStrategyKey>,
1733 > = IndexMap::new();
1734
1735 for key in quantities.keys() {
1736 keys_by_position
1737 .entry((key.0, key.1))
1738 .or_default()
1739 .push(*key);
1740 }
1741
1742 for (position_key, keys) in keys_by_position {
1743 let expected_qty = expected_quantities.get(&position_key).copied();
1744 let matches_report = if expected_qty.is_some_and(|quantity| quantity.is_zero()) {
1745 keys.iter().all(|key| {
1746 quantities
1747 .get(key)
1748 .copied()
1749 .flatten()
1750 .is_some_and(|quantity| quantity.is_zero())
1751 })
1752 } else if let (Some(expected_qty), [key]) = (expected_qty, keys.as_slice()) {
1753 let cache = self.cache.borrow();
1754 let positions = cache.positions_open(
1755 None,
1756 Some(&position_key.1),
1757 None,
1758 Some(&position_key.0),
1759 None,
1760 );
1761 let cache_is_unambiguous = positions.len() <= 1
1762 && positions.first().is_none_or(|position| {
1763 position.strategy_id == key.2
1764 && cache.oms_type(&position.id) == Some(OmsType::Netting)
1765 });
1766 cache_is_unambiguous
1767 && quantities
1768 .get(key)
1769 .copied()
1770 .flatten()
1771 .is_some_and(|quantity| quantity == expected_qty)
1772 } else {
1773 false
1774 };
1775
1776 if matches_report {
1777 continue;
1778 }
1779
1780 let venue_order_ids: Vec<VenueOrderId> = keys
1781 .iter()
1782 .filter_map(|key| group_ids.get(key))
1783 .flatten()
1784 .copied()
1785 .collect();
1786 log::error!(
1787 "Bounded reconciliation does not explain the reported position for {}; projecting {} historical order(s) without position or portfolio effects",
1788 position_key.1,
1789 venue_order_ids.len(),
1790 );
1791 order_only.extend(venue_order_ids);
1792 }
1793
1794 order_only
1795 }
1796
1797 pub fn check_inflight_orders(&mut self) -> InflightCheckResult {
1803 let mut result = InflightCheckResult::default();
1804 let now = dst::time::Instant::now();
1805 let threshold = Duration::from_millis(self.config.inflight_threshold_ms);
1806
1807 let mut to_check = Vec::new();
1808
1809 for (client_order_id, check) in &self.inflight_checks {
1810 if now
1811 .checked_duration_since(check.submitted_at)
1812 .is_some_and(|elapsed| elapsed > threshold)
1813 {
1814 to_check.push(*client_order_id);
1815 }
1816 }
1817
1818 for client_order_id in to_check {
1819 if self
1820 .config
1821 .filtered_client_order_ids
1822 .contains(&client_order_id)
1823 {
1824 self.clear_recon_tracking(&client_order_id, true);
1825 continue;
1826 }
1827
1828 if self.targeted_order_queries.contains(&client_order_id) {
1829 continue;
1830 }
1831
1832 if let Some(check) = self.inflight_checks.get_mut(&client_order_id) {
1833 if let Some(last_query_at) = check.last_query_at
1834 && now
1835 .checked_duration_since(last_query_at)
1836 .is_none_or(|elapsed| elapsed < threshold)
1837 {
1838 continue;
1839 }
1840
1841 check.retry_count += 1;
1842 check.last_query_at = Some(now);
1843 self.order_query_recency.mark(client_order_id);
1844 self.recon_check_retries
1845 .insert(client_order_id, check.retry_count);
1846
1847 if check.retry_count >= self.config.inflight_max_retries {
1848 let ts_now = self.clock.borrow().timestamp_ns();
1849
1850 if let Some(order) = self.get_order(client_order_id) {
1851 match order.status() {
1852 OrderStatus::Submitted => {
1853 if let Some(event) = create_reconciliation_rejected(
1855 &order,
1856 Some("INFLIGHT_TIMEOUT"),
1857 ts_now,
1858 ) {
1859 result.events.push(event);
1860 }
1861 }
1862 OrderStatus::PendingUpdate | OrderStatus::PendingCancel => {
1863 let event = OrderEventAny::Canceled(OrderCanceled::new(
1865 order.trader_id(),
1866 order.strategy_id(),
1867 order.instrument_id(),
1868 order.client_order_id(),
1869 UUID4::new(),
1870 ts_now,
1871 ts_now,
1872 true, order.venue_order_id(),
1874 order.account_id(),
1875 ));
1876 result.events.push(event);
1877 }
1878 _ => {
1879 }
1881 }
1882 }
1883 self.clear_recon_tracking(&client_order_id, true);
1885 } else if let Some(order) = self.get_order(client_order_id) {
1886 let ts_now = self.clock.borrow().timestamp_ns();
1888 let client_id = self.cache.borrow().client_id(&client_order_id).copied();
1889 let query = TradingCommand::QueryOrder(QueryOrder::new(
1890 order.trader_id(),
1891 client_id,
1892 order.strategy_id(),
1893 order.instrument_id(),
1894 order.client_order_id(),
1895 order.venue_order_id(),
1896 UUID4::new(),
1897 ts_now,
1898 None,
1899 None, ));
1901 result.queries.push(query);
1902 }
1903 }
1904 }
1905
1906 result
1907 }
1908
1909 pub(crate) fn validate_mass_status_order_sources(&self, mass_status: &ExecutionMassStatus) {
1913 let cache = self.cache.borrow();
1914 let mut checked_client_order_ids = IndexSet::new();
1915 let mut missing_origins: Vec<ClientOrderId> = Vec::new();
1916 let mut mismatched_origins: Vec<(ClientOrderId, ClientId)> = Vec::new();
1917
1918 let mut validate_report_source =
1919 |direct_client_order_id: Option<ClientOrderId>, venue_order_id: VenueOrderId| {
1920 let direct_client_order_id = direct_client_order_id
1921 .filter(|client_order_id| cache.order_exists(client_order_id));
1922 let indexed_client_order_id = cache
1923 .client_order_id(&venue_order_id)
1924 .copied()
1925 .filter(|client_order_id| cache.order_exists(client_order_id));
1926
1927 for client_order_id in [direct_client_order_id, indexed_client_order_id]
1928 .into_iter()
1929 .flatten()
1930 .filter(|client_order_id| checked_client_order_ids.insert(*client_order_id))
1931 {
1932 match cache.client_id(&client_order_id) {
1933 Some(cached_client_id) if *cached_client_id == mass_status.client_id => {}
1934 Some(cached_client_id) => {
1935 mismatched_origins.push((client_order_id, *cached_client_id));
1936 }
1937 None => missing_origins.push(client_order_id),
1938 }
1939 }
1940 };
1941
1942 for report in mass_status.order_reports().values() {
1943 validate_report_source(report.client_order_id, report.venue_order_id);
1944 }
1945
1946 for fills in mass_status.fill_reports().values() {
1947 for fill in fills {
1948 validate_report_source(fill.client_order_id, fill.venue_order_id);
1949 }
1950 }
1951
1952 if !missing_origins.is_empty() {
1953 let samples = missing_origins
1954 .iter()
1955 .take(5)
1956 .map(ToString::to_string)
1957 .collect::<Vec<_>>()
1958 .join(", ");
1959
1960 log::warn!(
1961 "Found {} cached order(s) without an execution client origin ({}): \
1962 continuing reconciliation against mass status client {} for compatibility \
1963 with existing cache data",
1964 missing_origins.len(),
1965 samples,
1966 mass_status.client_id,
1967 );
1968 }
1969
1970 if !mismatched_origins.is_empty() {
1971 let samples = mismatched_origins
1972 .iter()
1973 .take(5)
1974 .map(|(client_order_id, cached)| format!("{client_order_id} -> {cached}"))
1975 .collect::<Vec<_>>()
1976 .join(", ");
1977
1978 log::warn!(
1979 "Found {} cached order(s) with an execution client origin conflicting with \
1980 mass status client {} ({}): continuing reconciliation for compatibility; \
1981 this conflict will become a startup error in a future release, verify cached \
1982 order ownership and execution client configuration",
1983 mismatched_origins.len(),
1984 mass_status.client_id,
1985 samples,
1986 );
1987 }
1988 }
1989
1990 fn filtered_open_orders_for_reconciliation(&self) -> Vec<OrderAny> {
1991 {
1992 let cache = self.cache.borrow();
1993 let mut orders = cache.orders_open(None, None, None, None, None);
1994 orders.extend(cache.orders_inflight(None, None, None, None, None));
1995 let mut seen_client_order_ids = IndexSet::new();
1996 orders.retain(|order| seen_client_order_ids.insert(order.client_order_id()));
1997
1998 if self.config.reconciliation_instrument_ids.is_empty() {
1999 orders.iter().map(|o| (*o).clone()).collect()
2000 } else {
2001 orders
2002 .iter()
2003 .filter(|o| {
2004 self.config
2005 .reconciliation_instrument_ids
2006 .contains(&o.instrument_id())
2007 })
2008 .map(|o| (*o).clone())
2009 .collect()
2010 }
2011 }
2012 }
2013
2014 fn open_position_keys_for_reconciliation(&self) -> IndexSet<InstrumentAccountKey> {
2015 let cache = self.cache.borrow();
2016 let positions = cache.positions_open(None, None, None, None, None);
2017 let mut position_keys = IndexSet::new();
2018
2019 for position in positions {
2020 if !self.should_reconcile_instrument(&position.instrument_id) {
2021 continue;
2022 }
2023
2024 position_keys.insert((position.instrument_id, position.account_id));
2025 }
2026
2027 position_keys
2028 }
2029
2030 pub(crate) fn prepare_open_order_report_check(
2032 &mut self,
2033 command_id: UUID4,
2034 clients: &[&dyn ExecutionClient],
2035 ) -> OpenOrderReportCheck {
2036 let filtered_orders = self.filtered_open_orders_for_reconciliation();
2037 let active_order_ids: IndexSet<ClientOrderId> =
2038 filtered_orders.iter().map(Order::client_order_id).collect();
2039 self.missing_order_coverage_warnings
2040 .retain(|client_order_id| active_order_ids.contains(client_order_id));
2041 self.open_check_lookback_warnings
2042 .retain(|client_order_id| active_order_ids.contains(client_order_id));
2043 self.unresolved_order_coverage
2044 .retain(|client_order_id| active_order_ids.contains(client_order_id));
2045
2046 let mut client_coverage = IndexMap::new();
2047
2048 for order in &filtered_orders {
2049 let client_order_id = order.client_order_id();
2050 let coverage = self.resolve_order_report_client_coverage(order, clients);
2051
2052 match &coverage {
2053 ReportClientCoverage::Resolved(_) => {
2054 if self
2055 .unresolved_order_coverage
2056 .shift_remove(&client_order_id)
2057 {
2058 self.missing_order_coverage_warnings
2059 .shift_remove(&client_order_id);
2060 }
2061 }
2062 ReportClientCoverage::Unavailable(_) | ReportClientCoverage::Unresolved => {
2063 self.unresolved_order_coverage.insert(client_order_id);
2064 }
2065 }
2066
2067 client_coverage.insert(client_order_id, coverage);
2068 }
2069
2070 log::debug!(
2071 "Found {} order{} open in cache",
2072 filtered_orders.len(),
2073 if filtered_orders.len() == 1 { "" } else { "s" }
2074 );
2075
2076 let ts_now = self.clock.borrow().timestamp_ns();
2077 let start = self.config.open_check_lookback_mins.map(|mins| {
2078 let lookback_ns = mins_to_nanos(mins);
2079 ts_now.saturating_sub_ns(lookback_ns)
2080 });
2081
2082 let mut command = GenerateOrderStatusReports::new(
2083 command_id,
2084 ts_now,
2085 self.config.open_check_open_only,
2086 None,
2087 start,
2088 None,
2089 None,
2090 None,
2091 );
2092 command.log_receipt_level = LogLevel::Debug;
2093
2094 OpenOrderReportCheck {
2095 command,
2096 filtered_orders,
2097 client_coverage,
2098 start,
2099 }
2100 }
2101
2102 fn resolve_order_report_client_coverage(
2103 &self,
2104 order: &OrderAny,
2105 clients: &[&dyn ExecutionClient],
2106 ) -> ReportClientCoverage {
2107 if let Some(client_id) = self.cache.borrow().client_id(&order.client_order_id()) {
2108 return ReportClientCoverage::Resolved(IndexSet::from([*client_id]));
2109 }
2110
2111 if let Some(account_id) = order.account_id() {
2112 let account_clients = clients
2113 .iter()
2114 .filter(|client| client.account_id() == account_id)
2115 .map(|client| client.client_id())
2116 .collect::<IndexSet<_>>();
2117
2118 if !account_clients.is_empty() {
2119 return ReportClientCoverage::Resolved(account_clients);
2120 }
2121 }
2122
2123 let venue_clients = clients
2124 .iter()
2125 .filter(|client| client.handles_order_venue(order.instrument_id().venue))
2126 .map(|client| client.client_id())
2127 .collect::<IndexSet<_>>();
2128
2129 if venue_clients.is_empty() {
2130 ReportClientCoverage::Unresolved
2131 } else {
2132 ReportClientCoverage::Resolved(venue_clients)
2133 }
2134 }
2135
2136 pub fn check_open_order_queries(&mut self) -> Vec<TradingCommand> {
2138 self.check_open_order_queries_for_clients(None)
2139 }
2140
2141 pub(crate) fn check_open_order_queries_for_clients(
2142 &mut self,
2143 client_ids: Option<&IndexSet<ClientId>>,
2144 ) -> Vec<TradingCommand> {
2145 let now = dst::time::Instant::now();
2146 let query_delay = Duration::from_millis(u64::from(self.config.single_order_query_delay_ms));
2147 let query_limit = self.config.max_single_order_queries_per_cycle as usize;
2148
2149 if query_limit == 0 {
2150 return Vec::new();
2151 }
2152
2153 let mut filtered_orders = self.filtered_open_orders_for_reconciliation();
2154 filtered_orders.sort_by_key(|order| {
2155 let client_order_id = order.client_order_id();
2156 (
2157 self.order_query_recency.last_marked(&client_order_id),
2158 client_order_id,
2159 )
2160 });
2161
2162 let mut queries = Vec::new();
2163
2164 for order in filtered_orders {
2165 if queries.len() >= query_limit {
2166 break;
2167 }
2168
2169 let client_order_id = order.client_order_id();
2170 let client_id = self.cache.borrow().client_id(&client_order_id).copied();
2171
2172 if let Some(client_ids) = client_ids
2173 && !client_id.is_some_and(|client_id| client_ids.contains(&client_id))
2174 {
2175 continue;
2176 }
2177
2178 if self
2179 .config
2180 .filtered_client_order_ids
2181 .contains(&client_order_id)
2182 {
2183 continue;
2184 }
2185
2186 let threshold = Duration::from_nanos(self.config.open_check_threshold_ns);
2187 if let Some(elapsed) = self.order_local_activity.elapsed_at(&client_order_id, now)
2188 && elapsed < threshold
2189 {
2190 let elapsed_ms = elapsed.as_millis();
2191 let threshold_ms = threshold.as_millis();
2192 log::debug!(
2193 "Deferring open order query for {client_order_id}: recent local activity \
2194 ({elapsed_ms}ms < threshold={threshold_ms}ms)",
2195 );
2196 continue;
2197 }
2198
2199 if self
2200 .order_query_recency
2201 .within_at(&client_order_id, now, query_delay)
2202 {
2203 continue;
2204 }
2205
2206 self.order_query_recency.mark(client_order_id);
2207 let ts_now = self.clock.borrow().timestamp_ns();
2208
2209 let cmd = TradingCommand::QueryOrder(QueryOrder::new(
2210 order.trader_id(),
2211 client_id,
2212 order.strategy_id(),
2213 order.instrument_id(),
2214 client_order_id,
2215 order.venue_order_id(),
2216 UUID4::new(),
2217 ts_now,
2218 None,
2219 None,
2220 ));
2221 queries.push(cmd);
2222 }
2223
2224 queries
2225 }
2226
2227 pub async fn check_open_orders(
2237 &mut self,
2238 clients: &[&dyn ExecutionClient],
2239 ) -> Vec<OrderEventAny> {
2240 log::debug!("Checking order consistency between cached-state and venues");
2241
2242 let check = self.prepare_open_order_report_check(UUID4::new(), clients);
2243 let mut all_reports = Vec::new();
2244 let mut queried_clients = IndexSet::new();
2245 let mut failed_clients = IndexSet::new();
2246
2247 for client in clients {
2248 let client_id = client.client_id();
2249 queried_clients.insert(client_id);
2250
2251 match client.generate_order_status_reports(&check.command).await {
2252 Ok(reports) => {
2253 all_reports.extend(
2254 reports
2255 .into_iter()
2256 .map(|report| SourcedOrderStatusReport { client_id, report }),
2257 );
2258 }
2259 Err(e) => {
2260 failed_clients.insert(client_id);
2261 log::warn!(
2262 "Failed to query order reports from {}: {e}",
2263 client.client_id()
2264 );
2265 }
2266 }
2267 }
2268
2269 let result = self.reconcile_open_order_reports(
2270 &check,
2271 all_reports,
2272 &queried_clients,
2273 &failed_clients,
2274 clients,
2275 );
2276 let mut events = result.events;
2277
2278 if !result.targeted_queries.is_empty() {
2279 let query_delay =
2280 Duration::from_millis(u64::from(self.config.single_order_query_delay_ms));
2281 let query_results =
2282 request_targeted_order_reports(clients, result.targeted_queries, query_delay).await;
2283 events.extend(self.reconcile_targeted_order_reports(query_results, clients));
2284 }
2285
2286 events
2287 }
2288
2289 pub(crate) fn reconcile_open_order_reports(
2291 &mut self,
2292 check: &OpenOrderReportCheck,
2293 all_reports: Vec<SourcedOrderStatusReport>,
2294 queried_clients: &IndexSet<ClientId>,
2295 failed_clients: &IndexSet<ClientId>,
2296 clients: &[&dyn ExecutionClient],
2297 ) -> OpenOrderReconciliationResult {
2298 let mut venue_reported_ids = IndexSet::new();
2299
2300 for sourced in &all_reports {
2301 let report = &sourced.report;
2302 if let Some(client_order_id) = &report.client_order_id {
2303 venue_reported_ids.insert(*client_order_id);
2304 self.missing_order_coverage_warnings
2305 .shift_remove(client_order_id);
2306 self.open_check_lookback_warnings
2307 .shift_remove(client_order_id);
2308 self.recon_check_retries.shift_remove(client_order_id);
2312 } else {
2313 let mapped_client_order_id = self
2314 .cache
2315 .borrow()
2316 .client_order_id(&report.venue_order_id)
2317 .copied();
2318
2319 if let Some(client_order_id) = mapped_client_order_id {
2323 venue_reported_ids.insert(client_order_id);
2324 self.missing_order_coverage_warnings
2325 .shift_remove(&client_order_id);
2326 self.open_check_lookback_warnings
2327 .shift_remove(&client_order_id);
2328 self.recon_check_retries.shift_remove(&client_order_id);
2329 }
2330 }
2331 }
2332
2333 let mut events = Vec::new();
2334 let mut targeted_candidates = Vec::new();
2335
2336 for sourced in all_reports {
2337 let report = sourced.report;
2338 if let Some(client_order_id) = &report.client_order_id
2339 && let Some(order) = self.get_order(*client_order_id)
2340 {
2341 let threshold = Duration::from_nanos(self.config.open_check_threshold_ns);
2343 if let Some(elapsed) = self.order_local_activity.elapsed(client_order_id)
2344 && elapsed < threshold
2345 {
2346 let elapsed_ms = elapsed.as_millis();
2347 let threshold_ms = threshold.as_millis();
2348 log::debug!(
2349 "Deferring reconciliation for {client_order_id}: recent local activity ({elapsed_ms}ms < threshold={threshold_ms}ms)",
2350 );
2351 continue;
2352 }
2353
2354 let instrument = self.get_instrument(&report.instrument_id);
2355 let commission_client = clients
2356 .iter()
2357 .find(|client| client.client_id() == sourced.client_id)
2358 .copied();
2359
2360 match self.reconcile_order_report(
2361 &order,
2362 &report,
2363 instrument.as_ref(),
2364 commission_client,
2365 ) {
2366 Ok(Some(event)) => events.push(event),
2367 Ok(None) => {}
2368 Err(e) => log::error!(
2369 "Deferring reconciliation for {client_order_id}: venue commission calculation failed: {e}"
2370 ),
2371 }
2372 }
2373 }
2374
2375 if self.config.open_check_open_only {
2380 let cached_ids: IndexSet<ClientOrderId> = check
2381 .filtered_orders
2382 .iter()
2383 .map(Order::client_order_id)
2384 .collect();
2385 let missing_at_venue: IndexSet<ClientOrderId> = cached_ids
2386 .difference(&venue_reported_ids)
2387 .copied()
2388 .collect();
2389
2390 if !missing_at_venue.is_empty() {
2391 log::debug!(
2392 "{} cached open order{} not present in venue current response",
2393 missing_at_venue.len(),
2394 if missing_at_venue.len() == 1 {
2395 " is"
2396 } else {
2397 "s are"
2398 },
2399 );
2400
2401 for client_order_id in missing_at_venue {
2402 log::debug!("Cached open order missing from venue response: {client_order_id}");
2403 }
2404 }
2405 } else {
2406 let candidates: Vec<&OrderAny> = if let Some(cutoff) = check.start {
2407 let mut candidates = Vec::new();
2408
2409 for order in &check.filtered_orders {
2410 let client_order_id = order.client_order_id();
2411 if order.ts_last() >= cutoff {
2412 self.open_check_lookback_warnings
2413 .shift_remove(&client_order_id);
2414 candidates.push(order);
2415 } else if !venue_reported_ids.contains(&client_order_id)
2416 && self.open_check_lookback_warnings.insert(client_order_id)
2417 {
2418 log::warn!(
2419 "Skipping missing-order reconciliation for {client_order_id}: its last update predates the configured open-check lookback window; absence from the bulk response cannot be treated as evidence and no targeted query will be issued from it"
2420 );
2421 }
2422 }
2423 candidates
2424 } else {
2425 check.filtered_orders.iter().collect()
2426 };
2427
2428 for order in candidates {
2429 let client_order_id = order.client_order_id();
2430 if venue_reported_ids.contains(&client_order_id) {
2431 continue;
2432 }
2433
2434 let coverage = check
2435 .client_coverage
2436 .get(&client_order_id)
2437 .unwrap_or(&ReportClientCoverage::Unresolved);
2438
2439 let ReportClientCoverage::Resolved(responsible_clients) = coverage else {
2440 if self.missing_order_coverage_warnings.insert(client_order_id) {
2441 log::warn!(
2442 "Skipping order reconciliation for {client_order_id}: responsible execution client coverage is unresolved"
2443 );
2444 }
2445 continue;
2446 };
2447
2448 if responsible_clients.is_empty() {
2449 if self.missing_order_coverage_warnings.insert(client_order_id) {
2450 log::warn!(
2451 "Skipping order reconciliation for {client_order_id}: responsible execution client coverage is unresolved"
2452 );
2453 }
2454 continue;
2455 }
2456
2457 let missing_clients = responsible_clients
2458 .difference(queried_clients)
2459 .copied()
2460 .collect::<IndexSet<_>>();
2461
2462 if !missing_clients.is_empty() {
2463 if self.missing_order_coverage_warnings.insert(client_order_id) {
2464 log::warn!(
2465 "Skipping order reconciliation for {client_order_id}: responsible execution clients were not queried: {missing_clients:?}"
2466 );
2467 }
2468 continue;
2469 }
2470
2471 let failed_responsible_clients = responsible_clients
2472 .intersection(failed_clients)
2473 .copied()
2474 .collect::<IndexSet<_>>();
2475
2476 if !failed_responsible_clients.is_empty() {
2477 log::warn!(
2478 "Skipping order reconciliation for {client_order_id}: failed to query responsible execution clients: {failed_responsible_clients:?}"
2479 );
2480 continue;
2481 }
2482
2483 self.missing_order_coverage_warnings
2484 .shift_remove(&client_order_id);
2485 if let Some(order) = self.prepare_missing_order_query(client_order_id) {
2486 targeted_candidates.push((order, responsible_clients.clone()));
2487 }
2488 }
2489 }
2490
2491 targeted_candidates.sort_by_key(|(order, _)| {
2492 let client_order_id = order.client_order_id();
2493 (
2494 self.order_query_recency.last_marked(&client_order_id),
2495 client_order_id,
2496 )
2497 });
2498
2499 let query_limit = self.config.max_single_order_queries_per_cycle as usize;
2500 let mut planned_queries = 0usize;
2501 let mut cap_deferred_orders = 0usize;
2502 let mut targeted_queries = Vec::new();
2503
2504 for (order, responsible_clients) in targeted_candidates {
2505 let client_order_id = order.client_order_id();
2506
2507 let required_queries = responsible_clients.len();
2508 let exceeds_query_limit = planned_queries + required_queries > query_limit;
2509 let can_run_oversized_group = planned_queries == 0 && query_limit > 0;
2510 if required_queries == 0 || (exceeds_query_limit && !can_run_oversized_group) {
2511 cap_deferred_orders += 1;
2512 continue;
2513 }
2514
2515 if required_queries > query_limit {
2516 log::warn!(
2517 "Targeted order query for {client_order_id} requires {required_queries} responsible clients, exceeding the per-cycle limit {query_limit} to avoid indefinite deferral"
2518 );
2519 }
2520
2521 planned_queries += required_queries;
2522 self.order_query_recency.mark(client_order_id);
2523 self.targeted_order_queries.insert(client_order_id);
2524 targeted_queries.push(TargetedOrderQuery {
2525 client_order_id,
2526 responsible_clients,
2527 command: GenerateOrderStatusReport::new(
2528 UUID4::new(),
2529 self.clock.borrow().timestamp_ns(),
2530 Some(order.instrument_id()),
2531 Some(client_order_id),
2532 order.venue_order_id(),
2533 None,
2534 None,
2535 ),
2536 });
2537 }
2538
2539 if cap_deferred_orders > 0 {
2540 log::warn!(
2541 "Reached max single-order queries ({query_limit}) this cycle, deferring {cap_deferred_orders} order(s)"
2542 );
2543 }
2544
2545 OpenOrderReconciliationResult {
2546 events,
2547 targeted_queries,
2548 }
2549 }
2550
2551 pub(crate) fn reconcile_targeted_order_reports(
2552 &mut self,
2553 results: Vec<TargetedOrderReportResult>,
2554 clients: &[&dyn ExecutionClient],
2555 ) -> Vec<OrderEventAny> {
2556 let mut events = Vec::new();
2557
2558 for result in results {
2559 let client_order_id = result.client_order_id;
2560 self.targeted_order_queries.shift_remove(&client_order_id);
2561
2562 if let Some(report) = result.report {
2563 self.recon_check_retries.shift_remove(&client_order_id);
2564 self.missing_order_coverage_warnings
2565 .shift_remove(&client_order_id);
2566
2567 let Some(order) = self.get_order(client_order_id) else {
2568 continue;
2569 };
2570 let instrument = self.get_instrument(&report.instrument_id);
2571 let commission_client = result.client_id.and_then(|client_id| {
2572 clients
2573 .iter()
2574 .find(|client| client.client_id() == client_id)
2575 .copied()
2576 });
2577
2578 log::info!(
2579 color = LogColor::Blue as u8;
2580 "Found {client_order_id} via targeted order status query: {}",
2581 report.order_status,
2582 );
2583
2584 match self.reconcile_order_report(
2585 &order,
2586 &report,
2587 instrument.as_ref(),
2588 commission_client,
2589 ) {
2590 Ok(Some(event)) => events.push(event),
2591 Ok(None) => {}
2592 Err(e) => log::error!(
2593 "Deferring targeted reconciliation for {client_order_id}: venue commission calculation failed: {e}"
2594 ),
2595 }
2596 continue;
2597 }
2598
2599 if result.coverage_complete {
2600 events.extend(self.resolve_missing_order(client_order_id));
2601 } else {
2602 log::warn!(
2603 "Deferring missing-order resolution for {client_order_id}: targeted order status coverage was incomplete"
2604 );
2605 }
2606 }
2607
2608 events
2609 }
2610
2611 #[must_use]
2613 pub(crate) fn prepare_position_report_check(
2614 &self,
2615 command_id: UUID4,
2616 clients: &[&dyn ExecutionClient],
2617 ) -> PositionReportCheck {
2618 let position_keys = self.open_position_keys_for_reconciliation();
2619 let client_coverage = position_keys
2620 .iter()
2621 .map(|key| {
2622 (
2623 *key,
2624 Self::resolve_position_report_client_coverage(*key, clients),
2625 )
2626 })
2627 .collect();
2628 let mut activity_revisions = self.position_local_activity_revisions.clone();
2629 for key in &position_keys {
2630 activity_revisions
2631 .entry(*key)
2632 .or_insert_with(|| self.position_activity_revision(key));
2633 }
2634
2635 log::debug!(
2636 "Found {} unique instrument/account combination{} with open positions",
2637 position_keys.len(),
2638 if position_keys.len() == 1 { "" } else { "s" }
2639 );
2640
2641 let mut command = GeneratePositionStatusReports::new(
2642 command_id,
2643 self.clock.borrow().timestamp_ns(),
2644 None, None, None, None, None, );
2650 command.log_receipt_level = LogLevel::Debug;
2651
2652 PositionReportCheck {
2653 command,
2654 client_coverage,
2655 activity_revisions,
2656 }
2657 }
2658
2659 #[cfg(feature = "node")]
2660 pub(crate) fn prepare_position_fill_report_plan(
2661 &mut self,
2662 check: &mut PositionReportCheck,
2663 reports: &[PositionStatusReport],
2664 queried_clients: &IndexSet<ClientId>,
2665 failed_clients: &IndexSet<ClientId>,
2666 clients: &[&dyn ExecutionClient],
2667 ) -> PositionFillReportPlan {
2668 let mut venue_positions: IndexMap<InstrumentAccountKey, Vec<PositionStatusReport>> =
2669 IndexMap::new();
2670
2671 for report in reports {
2672 if self.should_reconcile_instrument(&report.instrument_id) {
2673 venue_positions
2674 .entry((report.instrument_id, report.account_id))
2675 .or_default()
2676 .push(report.clone());
2677 }
2678 }
2679
2680 let keys = check
2681 .client_coverage
2682 .keys()
2683 .copied()
2684 .chain(venue_positions.iter().filter_map(|(key, reports)| {
2685 reports
2686 .iter()
2687 .any(|report| report.signed_decimal_qty != Decimal::ZERO)
2688 .then_some(*key)
2689 }))
2690 .collect::<IndexSet<_>>();
2691 let active_keys = keys.clone();
2692 let query_end = self.clock.borrow().timestamp_ns();
2693 let lookback_ns = checked_mins_to_nanos(self.config.position_check_lookback_mins)
2694 .expect("position lookback validated at construction");
2695 let query_start = query_end.saturating_sub_ns(lookback_ns);
2696 let mut discrepancy_keys = IndexSet::new();
2697 let mut queries = Vec::new();
2698
2699 for key in keys {
2700 let coverage = check
2701 .client_coverage
2702 .entry(key)
2703 .or_insert_with(|| Self::resolve_position_report_client_coverage(key, clients));
2704 let prepared_revision = *check.activity_revisions.entry(key).or_default();
2705 let venue_reports = venue_positions
2706 .get(&key)
2707 .map(Vec::as_slice)
2708 .unwrap_or_default();
2709 let comparison = self.position_quantity_comparison(key, venue_reports);
2710 let tolerance = self.position_reconciliation_tolerance(key.1);
2711
2712 if comparison.quantities_match(tolerance) {
2713 self.position_reconciliation_states.shift_remove(&key);
2714 continue;
2715 }
2716 discrepancy_keys.insert(key);
2717
2718 if self.position_activity_revision(&key) > prepared_revision
2719 || self.position_local_activity.within(
2720 &key,
2721 Duration::from_nanos(self.config.position_check_threshold_ns),
2722 )
2723 {
2724 continue;
2725 }
2726
2727 let report_shape = comparison.report_shape();
2728 let retries = self
2729 .position_reconciliation_states
2730 .get(&key)
2731 .filter(|state| state.report_shape == report_shape)
2732 .map_or(0, |state| state.retries);
2733 if retries >= self.config.position_check_retries {
2734 continue;
2735 }
2736
2737 let ReportClientCoverage::Resolved(responsible_clients) = coverage else {
2738 log::warn!(
2739 "Skipping fill report query for {}/{}: responsible execution client coverage is unavailable",
2740 key.0,
2741 key.1,
2742 );
2743 continue;
2744 };
2745
2746 if responsible_clients.is_empty()
2747 || !responsible_clients.is_subset(queried_clients)
2748 || !responsible_clients.is_disjoint(failed_clients)
2749 {
2750 log::warn!(
2751 "Skipping fill report query for {}/{}: responsible position report coverage is incomplete",
2752 key.0,
2753 key.1,
2754 );
2755 continue;
2756 }
2757
2758 for client_id in responsible_clients.iter() {
2759 let mut command = GenerateFillReports::new(
2760 UUID4::new(),
2761 query_end,
2762 Some(key.0),
2763 None,
2764 Some(query_start),
2765 Some(query_end),
2766 None,
2767 Some(check.command.command_id),
2768 );
2769 command.log_receipt_level = LogLevel::Debug;
2770 queries.push(PositionFillReportQuery {
2771 key,
2772 client_id: *client_id,
2773 command,
2774 });
2775 }
2776 }
2777
2778 self.position_reconciliation_states
2779 .retain(|key, _| active_keys.contains(key));
2780
2781 PositionFillReportPlan {
2782 queries,
2783 discrepancy_keys,
2784 }
2785 }
2786
2787 #[cfg(feature = "node")]
2788 pub(crate) fn position_report_check_key_is_stable(
2789 &self,
2790 check: &PositionReportCheck,
2791 key: &InstrumentAccountKey,
2792 ) -> bool {
2793 check
2794 .activity_revisions
2795 .get(key)
2796 .is_some_and(|revision| self.position_activity_revision(key) == *revision)
2797 }
2798
2799 #[cfg(feature = "node")]
2800 pub(crate) fn prepare_position_fill_report(
2801 &self,
2802 report: &mut FillReport,
2803 venue_reports: &[PositionStatusReport],
2804 ) -> anyhow::Result<PositionFillReportPreparation> {
2805 let cache = self.cache.borrow();
2806 let venue_client_order_id = cache.client_order_id(&report.venue_order_id).copied();
2807 if let (Some(report_client_order_id), Some(venue_client_order_id)) =
2808 (report.client_order_id, venue_client_order_id)
2809 {
2810 anyhow::ensure!(
2811 report_client_order_id == venue_client_order_id,
2812 "fill {} client order ID {report_client_order_id} conflicts with venue order mapping {venue_client_order_id}",
2813 report.trade_id,
2814 );
2815 }
2816 let client_order_id = report.client_order_id.or(venue_client_order_id);
2817 let order = client_order_id.and_then(|id| cache.order(&id));
2818 if let Some(order) = &order {
2819 anyhow::ensure!(
2820 order.instrument_id() == report.instrument_id
2821 && order.order_side() == report.order_side
2822 && order
2823 .account_id()
2824 .is_none_or(|account_id| account_id == report.account_id)
2825 && order
2826 .venue_order_id()
2827 .is_none_or(|venue_order_id| venue_order_id == report.venue_order_id),
2828 "fill {} conflicts with cached order {}",
2829 report.trade_id,
2830 order.client_order_id(),
2831 );
2832 }
2833
2834 let hedge_context = report.venue_position_id.is_some()
2835 || venue_reports
2836 .iter()
2837 .any(|venue_report| venue_report.venue_position_id.is_some());
2838 let mapped_position_id = client_order_id
2839 .and_then(|client_order_id| cache.position_id(&client_order_id))
2840 .copied();
2841
2842 if hedge_context
2843 && let (Some(venue_position_id), Some(mapped_position_id)) =
2844 (report.venue_position_id, mapped_position_id)
2845 {
2846 anyhow::ensure!(
2847 venue_position_id == mapped_position_id,
2848 "fill {} position ID {venue_position_id} conflicts with cached order position {mapped_position_id}",
2849 report.trade_id,
2850 );
2851 }
2852
2853 if let Some(order) = order
2854 && Self::has_active_inferred_fill(&order)?
2855 {
2856 return Ok(PositionFillReportPreparation::InferredOverlap);
2857 }
2858
2859 if !hedge_context {
2860 return Ok(PositionFillReportPreparation::Ready);
2861 }
2862
2863 if report.venue_position_id.is_some() {
2864 return Ok(PositionFillReportPreparation::Ready);
2865 }
2866
2867 let Some(position_id) = mapped_position_id else {
2868 return Ok(PositionFillReportPreparation::Unattributed);
2869 };
2870 let position = cache.position(&position_id).ok_or_else(|| {
2871 anyhow::anyhow!(
2872 "fill {} maps to position {position_id}, which is not cached",
2873 report.trade_id,
2874 )
2875 })?;
2876
2877 anyhow::ensure!(
2878 position.account_id == report.account_id
2879 && position.instrument_id == report.instrument_id,
2880 "fill {} maps to position {position_id} with a different account or instrument",
2881 report.trade_id,
2882 );
2883 anyhow::ensure!(
2884 position.is_open(),
2885 "fill {} maps to non-open position {position_id}",
2886 report.trade_id,
2887 );
2888 anyhow::ensure!(
2889 !position.is_opposite_side(report.order_side) || report.last_qty <= position.quantity,
2890 "fill {} without a venue position ID would cross position {position_id}",
2891 report.trade_id,
2892 );
2893
2894 report.venue_position_id = Some(position_id);
2895 Ok(PositionFillReportPreparation::Ready)
2896 }
2897
2898 #[cfg(feature = "node")]
2899 fn has_active_inferred_fill(order: &OrderAny) -> anyhow::Result<bool> {
2900 let events = order.events();
2901 let trade_ids = order.trade_ids();
2902 let Some((first, remaining)) = events.split_first() else {
2903 return Ok(false);
2904 };
2905 let mut projected = OrderAny::from_events(vec![(*first).clone()]).map_err(|e| {
2906 anyhow::anyhow!(
2907 "cannot replay order {} for inferred fill detection: {e}",
2908 order.client_order_id(),
2909 )
2910 })?;
2911
2912 for event in remaining {
2913 projected.apply((*event).clone()).map_err(|e| {
2914 anyhow::anyhow!(
2915 "cannot replay order {} for inferred fill detection: {e}",
2916 order.client_order_id(),
2917 )
2918 })?;
2919 let OrderEventAny::Filled(fill) = event else {
2920 continue;
2921 };
2922
2923 if !fill.reconciliation || !trade_ids.contains(&&fill.trade_id) {
2924 continue;
2925 }
2926
2927 let external_position_id = PositionId::new(format!("{}-EXTERNAL", fill.instrument_id));
2928 let position_ids = [fill.position_id, Some(external_position_id)];
2929 let inferred = position_ids.into_iter().flatten().any(|position_id| {
2930 create_inferred_reconciliation_trade_id(
2931 fill.account_id,
2932 fill.instrument_id,
2933 fill.client_order_id,
2934 Some(fill.venue_order_id),
2935 fill.order_side,
2936 fill.order_type,
2937 projected.filled_qty(),
2938 fill.last_qty,
2939 fill.last_px,
2940 position_id,
2941 fill.ts_event,
2942 ) == fill.trade_id
2943 });
2944
2945 if inferred {
2946 return Ok(true);
2947 }
2948 }
2949
2950 Ok(false)
2951 }
2952
2953 #[cfg(feature = "node")]
2954 pub(crate) fn position_contains_fill_report(&self, report: &FillReport) -> bool {
2955 let cache = self.cache.borrow();
2956 let client_order_id = report
2957 .client_order_id
2958 .or_else(|| cache.client_order_id(&report.venue_order_id).copied());
2959 let positions = cache.positions(
2960 None,
2961 Some(&report.instrument_id),
2962 None,
2963 Some(&report.account_id),
2964 None,
2965 );
2966 let mut matched = false;
2967 let mut quantity_raw: QuantityRaw = 0;
2968 let mut commission_raw: MoneyRaw = 0;
2969
2970 for position in positions {
2971 if report
2972 .venue_position_id
2973 .is_some_and(|position_id| position.id != position_id)
2974 {
2975 continue;
2976 }
2977
2978 for replay_event in &position.replay_events {
2979 let PositionReplayEvent::Filled(fill) = replay_event else {
2980 continue;
2981 };
2982
2983 if fill.account_id != report.account_id
2984 || fill.instrument_id != report.instrument_id
2985 || fill.venue_order_id != report.venue_order_id
2986 || fill.trade_id != report.trade_id
2987 || fill.order_side != report.order_side
2988 || fill.last_px != report.last_px
2989 || fill.liquidity_side != report.liquidity_side
2990 || client_order_id.is_some_and(|id| fill.client_order_id != id)
2991 || report
2992 .venue_position_id
2993 .is_some_and(|id| fill.position_id != Some(id))
2994 {
2995 continue;
2996 }
2997
2998 let Some(fill_commission) = fill.commission else {
2999 return false;
3000 };
3001
3002 if fill_commission.currency != report.commission.currency {
3003 return false;
3004 }
3005 let Some(next_quantity_raw) = quantity_raw.checked_add(fill.last_qty.raw) else {
3006 return false;
3007 };
3008 let Some(next_commission_raw) = commission_raw.checked_add(fill_commission.raw)
3009 else {
3010 return false;
3011 };
3012 matched = true;
3013 quantity_raw = next_quantity_raw;
3014 commission_raw = next_commission_raw;
3015 }
3016 }
3017
3018 matched && quantity_raw == report.last_qty.raw && commission_raw == report.commission.raw
3019 }
3020
3021 fn resolve_position_report_client_coverage(
3022 key: InstrumentAccountKey,
3023 clients: &[&dyn ExecutionClient],
3024 ) -> ReportClientCoverage {
3025 let account_clients = clients
3026 .iter()
3027 .filter(|client| client.account_id() == key.1)
3028 .map(|client| client.client_id())
3029 .collect::<IndexSet<_>>();
3030
3031 if !account_clients.is_empty() {
3032 return if clients.iter().any(|client| {
3033 account_clients.contains(&client.client_id())
3034 && !client.provides_bulk_position_coverage(key.0)
3035 }) {
3036 ReportClientCoverage::Unavailable(account_clients)
3037 } else {
3038 ReportClientCoverage::Resolved(account_clients)
3039 };
3040 }
3041
3042 let venue_clients = clients
3043 .iter()
3044 .filter(|client| client.handles_order_venue(key.0.venue))
3045 .map(|client| client.client_id())
3046 .collect::<IndexSet<_>>();
3047
3048 if venue_clients.is_empty() {
3049 ReportClientCoverage::Unresolved
3050 } else if clients.iter().any(|client| {
3051 venue_clients.contains(&client.client_id())
3052 && !client.provides_bulk_position_coverage(key.0)
3053 }) {
3054 ReportClientCoverage::Unavailable(venue_clients)
3055 } else {
3056 ReportClientCoverage::Resolved(venue_clients)
3057 }
3058 }
3059
3060 pub async fn check_positions_consistency(
3070 &mut self,
3071 clients: &[&dyn ExecutionClient],
3072 ) -> Vec<OrderEventAny> {
3073 let check = self.prepare_position_report_check(UUID4::new(), clients);
3074 let mut reports = Vec::new();
3075 let mut queried_clients = IndexSet::new();
3076 let mut failed_clients = IndexSet::new();
3077
3078 for client in clients {
3079 let client_id = client.client_id();
3080 queried_clients.insert(client_id);
3081 self.set_position_reconciliation_tolerance(
3082 client.account_id(),
3083 client.position_reconciliation_tolerance(),
3084 );
3085
3086 match client
3087 .generate_position_status_reports(&check.command)
3088 .await
3089 {
3090 Ok(client_reports) => {
3091 reports.extend(client_reports);
3092 }
3093 Err(e) => {
3094 failed_clients.insert(client_id);
3095 log::warn!(
3096 "Failed to query position reports from {}: {e}",
3097 client.client_id()
3098 );
3099 }
3100 }
3101 }
3102
3103 self.reconcile_position_reports(&check, reports, &queried_clients, &failed_clients)
3104 }
3105
3106 #[must_use]
3108 pub(crate) fn reconcile_position_reports(
3109 &mut self,
3110 check: &PositionReportCheck,
3111 reports: Vec<PositionStatusReport>,
3112 queried_clients: &IndexSet<ClientId>,
3113 failed_clients: &IndexSet<ClientId>,
3114 ) -> Vec<OrderEventAny> {
3115 log::debug!("Checking position consistency between cached-state and venues");
3116
3117 let mut venue_positions: IndexMap<InstrumentAccountKey, Vec<PositionStatusReport>> =
3118 IndexMap::new();
3119
3120 for report in reports {
3121 if !self.should_reconcile_instrument(&report.instrument_id) {
3122 continue;
3123 }
3124
3125 venue_positions
3126 .entry((report.instrument_id, report.account_id))
3127 .or_default()
3128 .push(report);
3129 }
3130
3131 let mut events = Vec::new();
3132
3133 for key in check.client_coverage.keys() {
3134 let prepared_revision = check
3135 .activity_revisions
3136 .get(key)
3137 .copied()
3138 .unwrap_or_default();
3139
3140 if self.position_activity_revision(key) > prepared_revision {
3141 log::debug!(
3142 "Deferring position reconciliation for {}/{}: local activity recorded during report request",
3143 key.0,
3144 key.1,
3145 );
3146 continue;
3147 }
3148
3149 let venue_reports = venue_positions
3150 .get(key)
3151 .map(Vec::as_slice)
3152 .unwrap_or_default();
3153
3154 if venue_reports.is_empty() {
3155 match check.client_coverage.get(key) {
3156 Some(ReportClientCoverage::Resolved(responsible_clients))
3157 if !responsible_clients.is_empty()
3158 && responsible_clients.is_subset(queried_clients)
3159 && responsible_clients.is_disjoint(failed_clients) => {}
3160 Some(ReportClientCoverage::Resolved(responsible_clients))
3161 if responsible_clients.is_empty() =>
3162 {
3163 log::warn!(
3164 "Skipping position reconciliation for {}/{}: responsible execution client coverage is unresolved",
3165 key.0,
3166 key.1,
3167 );
3168 continue;
3169 }
3170 Some(ReportClientCoverage::Resolved(responsible_clients))
3171 if !responsible_clients.is_subset(queried_clients) =>
3172 {
3173 log::warn!(
3174 "Skipping position reconciliation for {}/{}: responsible execution clients were not all queried",
3175 key.0,
3176 key.1,
3177 );
3178 continue;
3179 }
3180 Some(ReportClientCoverage::Resolved(responsible_clients)) => {
3181 let failed_responsible_clients = responsible_clients
3182 .intersection(failed_clients)
3183 .copied()
3184 .collect::<IndexSet<_>>();
3185 log::warn!(
3186 "Skipping position reconciliation for {}/{}: failed to query responsible execution clients: {failed_responsible_clients:?}",
3187 key.0,
3188 key.1,
3189 );
3190 continue;
3191 }
3192 Some(ReportClientCoverage::Unavailable(responsible_clients)) => {
3193 log::debug!(
3194 "Skipping position reconciliation for {}/{}: complete bulk position coverage is unavailable from responsible execution clients: {responsible_clients:?}",
3195 key.0,
3196 key.1,
3197 );
3198 continue;
3199 }
3200 Some(ReportClientCoverage::Unresolved) | None => {
3201 log::warn!(
3202 "Skipping position reconciliation for {}/{}: responsible execution client coverage is unresolved",
3203 key.0,
3204 key.1,
3205 );
3206 continue;
3207 }
3208 }
3209 }
3210
3211 if let Some(discrepancy_events) = self.check_position_discrepancy(*key, venue_reports) {
3212 events.extend(discrepancy_events);
3213 }
3214 }
3215
3216 let current_position_keys = self.open_position_keys_for_reconciliation();
3217
3218 for (key, venue_reports) in &venue_positions {
3219 if check.client_coverage.contains_key(key)
3220 || venue_reports
3221 .iter()
3222 .all(|report| report.signed_decimal_qty == Decimal::ZERO)
3223 {
3224 continue;
3225 }
3226
3227 if current_position_keys.contains(key) {
3228 log::debug!(
3229 "Deferring position reconciliation for {}/{}: position opened after client coverage was recorded",
3230 key.0,
3231 key.1,
3232 );
3233 continue;
3234 }
3235
3236 if let Some(discrepancy_events) = self.check_position_discrepancy(*key, venue_reports) {
3237 events.extend(discrepancy_events);
3238 }
3239 }
3240
3241 let active_keys: IndexSet<InstrumentAccountKey> = current_position_keys
3244 .into_iter()
3245 .chain(
3246 venue_positions
3247 .iter()
3248 .filter(|(_, reports)| {
3249 reports
3250 .iter()
3251 .any(|report| report.signed_decimal_qty != Decimal::ZERO)
3252 })
3253 .map(|(k, _)| *k),
3254 )
3255 .collect();
3256 self.position_reconciliation_states
3257 .retain(|k, _| active_keys.contains(k));
3258
3259 events
3260 }
3261
3262 fn positions_avg_px(cached_positions: &[Position]) -> Option<Decimal> {
3263 let mut total_value = Decimal::ZERO;
3264 let mut total_qty = Decimal::ZERO;
3265
3266 for position in cached_positions {
3267 let qty = position.signed_decimal_qty().abs();
3268 if position.avg_px_open > 0.0
3269 && qty > Decimal::ZERO
3270 && let Ok(avg_px) = Decimal::from_str(&position.avg_px_open.to_string())
3271 {
3272 total_value += avg_px * qty;
3273 total_qty += qty;
3274 }
3275 }
3276
3277 if total_qty > Decimal::ZERO {
3278 Some(total_value / total_qty)
3279 } else {
3280 None
3281 }
3282 }
3283
3284 pub fn register_inflight(&mut self, client_order_id: ClientOrderId) {
3286 if self
3287 .config
3288 .filtered_client_order_ids
3289 .contains(&client_order_id)
3290 {
3291 return;
3292 }
3293
3294 self.inflight_checks.insert(
3295 client_order_id,
3296 InflightCheck {
3297 client_order_id,
3298 submitted_at: dst::time::Instant::now(),
3299 retry_count: 0,
3300 last_query_at: None,
3301 },
3302 );
3303 self.recon_check_retries.insert(client_order_id, 0);
3304 self.order_query_recency.remove(&client_order_id);
3305 self.order_local_activity.remove(&client_order_id);
3306 }
3307
3308 pub fn record_local_activity(&mut self, client_order_id: ClientOrderId) {
3315 self.order_local_activity.mark(client_order_id);
3316 }
3317
3318 pub fn clear_recon_tracking(&mut self, client_order_id: &ClientOrderId, drop_last_query: bool) {
3320 self.inflight_checks.shift_remove(client_order_id);
3321 self.recon_check_retries.shift_remove(client_order_id);
3322 self.missing_order_coverage_warnings
3323 .shift_remove(client_order_id);
3324 self.open_check_lookback_warnings
3325 .shift_remove(client_order_id);
3326 self.unresolved_order_coverage.shift_remove(client_order_id);
3327 self.targeted_order_queries.shift_remove(client_order_id);
3328
3329 if drop_last_query {
3330 self.order_query_recency.remove(client_order_id);
3331 }
3332 self.order_local_activity.remove(client_order_id);
3333 }
3334
3335 #[cfg(feature = "node")]
3336 pub(crate) fn remove_targeted_order_queries(&mut self, client_order_ids: &[ClientOrderId]) {
3337 for client_order_id in client_order_ids {
3338 self.targeted_order_queries.shift_remove(client_order_id);
3339 }
3340 }
3341
3342 #[must_use]
3344 pub fn get_external_order_claim(&self, instrument_id: &InstrumentId) -> Option<StrategyId> {
3345 self.external_order_claims.get(instrument_id).copied()
3346 }
3347
3348 #[must_use]
3350 #[cfg(feature = "node")]
3351 pub(crate) fn get_external_order_claims_for_strategy(
3352 &self,
3353 strategy_id: StrategyId,
3354 ) -> HashSet<InstrumentId> {
3355 self.external_order_claims
3356 .iter()
3357 .filter_map(|(instrument_id, owner)| (*owner == strategy_id).then_some(*instrument_id))
3358 .collect()
3359 }
3360
3361 pub fn claim_external_orders(
3367 &mut self,
3368 instrument_id: InstrumentId,
3369 strategy_id: StrategyId,
3370 ) -> anyhow::Result<()> {
3371 if let Some(existing) = self.external_order_claims.get(&instrument_id) {
3372 anyhow::bail!("External order claim for {instrument_id} already exists for {existing}");
3373 }
3374
3375 self.external_order_claims
3376 .insert(instrument_id, strategy_id);
3377 Ok(())
3378 }
3379
3380 #[cfg(feature = "node")]
3381 pub(crate) fn register_external_order_claims(
3382 &mut self,
3383 strategy_id: StrategyId,
3384 instrument_ids: &HashSet<InstrumentId>,
3385 ) {
3386 self.external_order_claims.extend(
3387 instrument_ids
3388 .iter()
3389 .map(|instrument_id| (*instrument_id, strategy_id)),
3390 );
3391 }
3392
3393 #[cfg(feature = "node")]
3399 pub(crate) fn deregister_external_order_claims(&mut self, strategy_id: StrategyId) {
3400 self.external_order_claims
3401 .retain(|_, owner| *owner != strategy_id);
3402 }
3403
3404 pub fn record_position_activity(&mut self, instrument_id: InstrumentId, account_id: AccountId) {
3418 let key = (instrument_id, account_id);
3419 self.position_local_activity.mark(key);
3420 let revision = self
3421 .position_local_activity_revisions
3422 .entry(key)
3423 .or_default();
3424 *revision = revision.saturating_add(1);
3425 }
3426
3427 pub(crate) fn position_activity_revision(&self, key: &InstrumentAccountKey) -> u64 {
3428 self.position_local_activity_revisions
3429 .get(key)
3430 .copied()
3431 .unwrap_or_default()
3432 }
3433
3434 #[must_use]
3437 pub fn position_recon_retry_count(&self, key: &InstrumentAccountKey) -> u32 {
3438 self.position_reconciliation_states
3439 .get(key)
3440 .map_or(0, |state| state.retries)
3441 }
3442
3443 #[must_use]
3446 pub fn recon_check_retry_count(&self, client_order_id: &ClientOrderId) -> u32 {
3447 self.recon_check_retries
3448 .get(client_order_id)
3449 .copied()
3450 .unwrap_or(0)
3451 }
3452
3453 pub fn observe_order_event(&mut self, event: &OrderEventAny) {
3463 match event {
3464 OrderEventAny::Filled(fill) => {
3465 self.record_position_activity(fill.instrument_id, fill.account_id);
3466 }
3467 OrderEventAny::Accepted(_)
3468 | OrderEventAny::Rejected(_)
3469 | OrderEventAny::Canceled(_)
3470 | OrderEventAny::Expired(_)
3471 | OrderEventAny::Denied(_)
3472 | OrderEventAny::Updated(_)
3473 | OrderEventAny::ModifyRejected(_)
3474 | OrderEventAny::CancelRejected(_) => {
3475 self.clear_recon_tracking(&event.client_order_id(), true);
3476 }
3477 _ => {}
3478 }
3479
3480 self.record_local_activity(event.client_order_id());
3481 }
3482
3483 pub fn observe_execution_report(&mut self, report: &ExecutionReport) {
3495 match report {
3496 ExecutionReport::Order(order_report) => {
3497 self.observe_order_status_report(order_report);
3498 }
3499 ExecutionReport::Fill(fill_report) => {
3500 let client_order_id = fill_report.client_order_id.or_else(|| {
3501 self.cache
3502 .borrow()
3503 .client_order_id(&fill_report.venue_order_id)
3504 .copied()
3505 });
3506
3507 if let Some(coid) = client_order_id {
3508 self.record_local_activity(coid);
3509 }
3510 self.record_position_activity(fill_report.instrument_id, fill_report.account_id);
3511 }
3512 ExecutionReport::OrderWithFills(order_report, fills) => {
3513 self.observe_order_status_report(order_report);
3514
3515 for fill_report in fills {
3516 self.record_position_activity(
3517 fill_report.instrument_id,
3518 fill_report.account_id,
3519 );
3520 }
3521 }
3522 ExecutionReport::Position(position_report) => {
3523 self.record_position_activity(
3524 position_report.instrument_id,
3525 position_report.account_id,
3526 );
3527 }
3528 ExecutionReport::MassStatus(_) => {
3529 }
3531 }
3532 }
3533
3534 fn observe_order_status_report(&mut self, report: &OrderStatusReport) {
3535 let Some(client_order_id) = report.client_order_id else {
3536 return;
3537 };
3538
3539 let accepted_during_pending_command = report.order_status == OrderStatus::Accepted
3540 && self.get_order(client_order_id).is_some_and(|order| {
3541 matches!(
3542 order.status(),
3543 OrderStatus::PendingUpdate | OrderStatus::PendingCancel
3544 )
3545 });
3546
3547 if !matches!(
3548 report.order_status,
3549 OrderStatus::PendingUpdate | OrderStatus::PendingCancel
3550 ) && !accepted_during_pending_command
3551 {
3552 self.clear_recon_tracking(&client_order_id, report.order_status.is_closed());
3553 }
3554
3555 self.record_local_activity(client_order_id);
3559 }
3560
3561 #[must_use]
3563 pub fn is_fill_recently_processed(
3564 &self,
3565 account_id: AccountId,
3566 instrument_id: InstrumentId,
3567 trade_id: TradeId,
3568 ) -> bool {
3569 self.recent_fills_cache
3570 .contains_key(&(account_id, instrument_id, trade_id))
3571 }
3572
3573 pub fn mark_fill_processed(
3575 &mut self,
3576 account_id: AccountId,
3577 instrument_id: InstrumentId,
3578 trade_id: TradeId,
3579 ) {
3580 self.recent_fills_cache
3581 .mark((account_id, instrument_id, trade_id));
3582 }
3583
3584 pub fn commit_recent_fill_if_applied(&mut self, fill: &OrderFilled) {
3586 let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
3587 if self.is_fill_applied(fill, fill_key) {
3588 self.mark_fill_processed(fill_key.0, fill_key.1, fill_key.2);
3589 }
3590 }
3591
3592 pub fn prune_recent_fills_cache(&mut self, ttl_secs: f64) {
3596 let ttl = match Duration::try_from_secs_f64(ttl_secs) {
3604 Ok(ttl) => ttl,
3605 Err(_) if ttl_secs > 0.0 => Duration::MAX,
3606 Err(_) => Duration::ZERO,
3607 };
3608
3609 self.recent_fills_cache.prune_older_than(ttl);
3610 }
3611
3612 pub fn prune_processed_fills(&mut self) {
3617 let Some(lookback_mins) = self.config.lookback_mins else {
3618 return;
3619 };
3620
3621 let ttl = Duration::from_mins(lookback_mins).max(Duration::from_mins(1));
3622 self.processed_fills.prune_older_than(ttl);
3623 }
3624
3625 pub fn prune_order_local_activity(&mut self) {
3627 self.order_local_activity
3628 .prune_older_than(Duration::from_nanos(self.config.open_check_threshold_ns));
3629 }
3630
3631 pub fn purge_closed_orders(&mut self) {
3633 let Some(buffer_mins) = self.config.purge_closed_orders_buffer_mins else {
3634 return;
3635 };
3636
3637 let ts_now = self.clock.borrow().timestamp_ns();
3638 let buffer_secs = mins_to_secs(u64::from(buffer_mins));
3639
3640 self.cache
3641 .borrow_mut()
3642 .purge_closed_orders(ts_now, buffer_secs);
3643 }
3644
3645 pub fn purge_closed_positions(&mut self) {
3647 let Some(buffer_mins) = self.config.purge_closed_positions_buffer_mins else {
3648 return;
3649 };
3650
3651 let ts_now = self.clock.borrow().timestamp_ns();
3652 let buffer_secs = mins_to_secs(u64::from(buffer_mins));
3653
3654 self.cache
3655 .borrow_mut()
3656 .purge_closed_positions(ts_now, buffer_secs);
3657 }
3658
3659 pub fn purge_account_events(&mut self) {
3661 let Some(lookback_mins) = self.config.purge_account_events_lookback_mins else {
3662 return;
3663 };
3664
3665 let ts_now = self.clock.borrow().timestamp_ns();
3666 let lookback_secs = mins_to_secs(u64::from(lookback_mins));
3667
3668 self.cache
3669 .borrow_mut()
3670 .purge_account_events(ts_now, lookback_secs);
3671 }
3672
3673 fn get_order(&self, client_order_id: ClientOrderId) -> Option<OrderAny> {
3676 self.cache
3677 .borrow()
3678 .order(&client_order_id)
3679 .map(|o| o.clone())
3680 }
3681
3682 fn get_order_by_venue_order_id(&self, venue_order_id: VenueOrderId) -> Option<OrderAny> {
3683 let cache = self.cache.borrow();
3684 cache
3685 .client_order_id(&venue_order_id)
3686 .and_then(|client_order_id| cache.order(client_order_id).map(|o| o.clone()))
3687 }
3688
3689 fn get_instrument(&self, instrument_id: &InstrumentId) -> Option<InstrumentAny> {
3690 self.cache.borrow().instrument(instrument_id).cloned()
3691 }
3692
3693 fn should_skip_order_report(&self, report: &OrderStatusReport) -> bool {
3694 if let Some(client_order_id) = &report.client_order_id
3695 && self
3696 .config
3697 .filtered_client_order_ids
3698 .contains(client_order_id)
3699 {
3700 log::debug!(
3701 "Skipping order report {client_order_id}: in filtered_client_order_ids list"
3702 );
3703 return true;
3704 }
3705
3706 if !self.should_reconcile_instrument(&report.instrument_id) {
3707 log::debug!(
3708 "Skipping order report for {}: not in reconciliation_instrument_ids",
3709 report.instrument_id
3710 );
3711 return true;
3712 }
3713
3714 false
3715 }
3716
3717 fn should_reconcile_instrument(&self, instrument_id: &InstrumentId) -> bool {
3718 self.config.reconciliation_instrument_ids.is_empty()
3719 || self
3720 .config
3721 .reconciliation_instrument_ids
3722 .contains(instrument_id)
3723 }
3724
3725 fn prepare_missing_order_query(&mut self, client_order_id: ClientOrderId) -> Option<OrderAny> {
3726 let order = self.get_order(client_order_id)?;
3727
3728 if order.status().is_closed() {
3732 log::debug!(
3733 "Skipping missing-order resolution for {client_order_id}: already {}",
3734 order.status()
3735 );
3736 self.clear_recon_tracking(&client_order_id, true);
3737 return None;
3738 }
3739
3740 if self.order_local_activity.within(
3744 &client_order_id,
3745 Duration::from_nanos(self.config.open_check_threshold_ns),
3746 ) {
3747 return None;
3748 }
3749
3750 let retries = self.recon_check_retries.entry(client_order_id).or_insert(0);
3751 *retries = retries.saturating_add(1);
3752
3753 if *retries < self.config.open_check_missing_retries {
3754 log::debug!(
3755 "Order {} not found at venue, retry {}/{}",
3756 client_order_id,
3757 retries,
3758 self.config.open_check_missing_retries
3759 );
3760 return None;
3761 }
3762
3763 Some(order)
3764 }
3765
3766 fn resolve_missing_order(&mut self, client_order_id: ClientOrderId) -> Vec<OrderEventAny> {
3767 let mut events = Vec::new();
3768
3769 let Some(order) = self.get_order(client_order_id) else {
3770 return events;
3771 };
3772
3773 if order.status().is_closed() {
3774 log::debug!(
3775 "Skipping missing-order resolution for {client_order_id}: already {}",
3776 order.status()
3777 );
3778 self.clear_recon_tracking(&client_order_id, true);
3779 return events;
3780 }
3781
3782 if self.order_local_activity.within(
3783 &client_order_id,
3784 Duration::from_nanos(self.config.open_check_threshold_ns),
3785 ) {
3786 log::debug!(
3787 "Deferring missing-order resolution for {client_order_id}: recent local activity"
3788 );
3789 return events;
3790 }
3791
3792 let retries = self
3793 .recon_check_retries
3794 .get(&client_order_id)
3795 .copied()
3796 .unwrap_or_default();
3797 let ts_now = self.clock.borrow().timestamp_ns();
3798
3799 match order.status() {
3800 OrderStatus::Accepted | OrderStatus::Submitted => {
3801 log::warn!(
3802 "Order {client_order_id} not found at venue after {retries} retries and a targeted query, marking as REJECTED"
3803 );
3804
3805 if let Some(rejected) =
3806 create_reconciliation_rejected(&order, Some("NOT_FOUND_AT_VENUE"), ts_now)
3807 {
3808 events.push(rejected);
3809 }
3810 }
3811 OrderStatus::PartiallyFilled => {
3812 log::warn!(
3813 "Order {client_order_id} not found at venue after {retries} retries and a targeted query, marking as CANCELED"
3814 );
3815 events.push(OrderEventAny::Canceled(OrderCanceled::new(
3816 order.trader_id(),
3817 order.strategy_id(),
3818 order.instrument_id(),
3819 client_order_id,
3820 UUID4::new(),
3821 ts_now,
3822 ts_now,
3823 true,
3824 order.venue_order_id(),
3825 order.account_id(),
3826 )));
3827 }
3828 OrderStatus::PendingUpdate | OrderStatus::PendingCancel => {
3829 log::debug!(
3830 "Deferring resolution for {client_order_id}: still inflight as {}",
3831 order.status()
3832 );
3833 self.recon_check_retries.shift_remove(&client_order_id);
3842 if let Some(check) = self.inflight_checks.get_mut(&client_order_id) {
3843 check.retry_count = 0;
3844 check.last_query_at = Some(dst::time::Instant::now());
3845 }
3846 self.order_query_recency.mark(client_order_id);
3847 return events;
3848 }
3849 status => {
3850 log::warn!(
3851 "Skipping missing-order resolution for {client_order_id}: unexpected status {status}"
3852 );
3853 }
3854 }
3855
3856 self.clear_recon_tracking(&client_order_id, true);
3857 events
3858 }
3859
3860 fn position_quantity_comparison(
3861 &self,
3862 key: InstrumentAccountKey,
3863 venue_reports: &[PositionStatusReport],
3864 ) -> PositionQuantityComparison {
3865 let (instrument_id, account_id) = key;
3866 let cached_positions = {
3867 let cache = self.cache.borrow();
3868 cache
3869 .positions_open(None, Some(&instrument_id), None, Some(&account_id), None)
3870 .into_iter()
3871 .map(|position| (*position).clone())
3872 .collect::<Vec<_>>()
3873 };
3874 let (cached_signed_qty, cached_long_qty, cached_short_qty) = Self::position_qty_aggregates(
3875 cached_positions.iter().map(Position::signed_decimal_qty),
3876 );
3877 let (venue_signed_qty, venue_long_qty, venue_short_qty) = Self::position_qty_aggregates(
3878 venue_reports.iter().map(|report| report.signed_decimal_qty),
3879 );
3880 let nonflat_count = venue_reports
3881 .iter()
3882 .filter(|report| report.signed_decimal_qty != Decimal::ZERO)
3883 .count();
3884 let venue_report = venue_reports
3885 .iter()
3886 .find(|report| report.signed_decimal_qty != Decimal::ZERO)
3887 .or_else(|| venue_reports.last())
3888 .cloned();
3889 let venue_has_side_reports = venue_reports.iter().any(PositionStatusReport::is_long)
3890 && venue_reports.iter().any(PositionStatusReport::is_short);
3891
3892 PositionQuantityComparison {
3893 cached_positions,
3894 cached_signed_qty,
3895 cached_long_qty,
3896 cached_short_qty,
3897 venue_signed_qty,
3898 venue_long_qty,
3899 venue_short_qty,
3900 nonflat_count,
3901 venue_report,
3902 venue_has_side_reports,
3903 }
3904 }
3905
3906 fn check_position_discrepancy(
3907 &mut self,
3908 key: InstrumentAccountKey,
3909 venue_reports: &[PositionStatusReport],
3910 ) -> Option<Vec<OrderEventAny>> {
3911 let (instrument_id, account_id) = key;
3912 let comparison = self.position_quantity_comparison(key, venue_reports);
3913 let tolerance = self.position_reconciliation_tolerance(account_id);
3914 let quantities_match = comparison.quantities_match(tolerance);
3915 let report_shape = comparison.report_shape();
3916 let PositionQuantityComparison {
3917 cached_positions,
3918 cached_signed_qty,
3919 cached_long_qty,
3920 cached_short_qty,
3921 venue_signed_qty,
3922 venue_long_qty,
3923 venue_short_qty,
3924 venue_report,
3925 ..
3926 } = comparison;
3927
3928 if quantities_match {
3929 self.position_reconciliation_states.shift_remove(&key);
3930 return None;
3931 }
3932
3933 let ts_now = self.clock.borrow().timestamp_ns();
3934
3935 if self.position_local_activity.within(
3937 &key,
3938 Duration::from_nanos(self.config.position_check_threshold_ns),
3939 ) {
3940 log::debug!(
3941 "Skipping position reconciliation for {instrument_id}: recent activity within threshold"
3942 );
3943 return None;
3944 }
3945
3946 let retries = self
3947 .position_reconciliation_states
3948 .get(&key)
3949 .filter(|state| state.report_shape == report_shape)
3950 .map_or(0, |state| state.retries);
3951
3952 if retries >= self.config.position_check_retries {
3953 return None;
3954 }
3955
3956 if report_shape == PositionReportShape::MultiLeg {
3957 let new_retries = retries + 1;
3958 self.set_position_reconciliation_retries(key, report_shape, new_retries);
3959 log::warn!(
3960 "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})"
3961 );
3962
3963 if new_retries >= self.config.position_check_retries {
3964 log::error!(
3965 "Position discrepancy for {instrument_id}/{account_id} unresolved after {} attempts; no further reconciliation attempts will be made for the current report shape",
3966 self.config.position_check_retries,
3967 );
3968 }
3969 return None;
3970 }
3971
3972 log::warn!(
3973 "Position discrepancy detected for {instrument_id}: cached_signed_qty={cached_signed_qty}, venue_signed_qty={venue_signed_qty}"
3974 );
3975
3976 let Some(instrument) = self.cache.borrow().instrument(&instrument_id).cloned() else {
3977 log::debug!("Cannot reconcile position for {instrument_id}: instrument not in cache");
3978 let new_retries = retries + 1;
3979 self.set_position_reconciliation_retries(key, report_shape, new_retries);
3980 if new_retries >= self.config.position_check_retries {
3981 log::error!(
3982 "Position discrepancy for {instrument_id} unresolved after {} attempts \
3983 (cached_qty={cached_signed_qty}, venue_qty={venue_signed_qty}); \
3984 no further reconciliation attempts will be made for the current report shape",
3985 self.config.position_check_retries,
3986 );
3987 }
3988 return None;
3989 };
3990
3991 let cached_avg_px = Self::positions_avg_px(&cached_positions);
3992 let venue_avg_px = venue_report.as_ref().and_then(|r| r.avg_px_open);
3993
3994 let crosses_zero = (cached_signed_qty > Decimal::ZERO && venue_signed_qty < Decimal::ZERO)
3995 || (cached_signed_qty < Decimal::ZERO && venue_signed_qty > Decimal::ZERO);
3996
3997 let result = if crosses_zero {
3998 let venue_ts_last = venue_report.as_ref().map_or(ts_now, |r| r.ts_last);
3999 let venue_position_id = venue_report
4000 .as_ref()
4001 .and_then(|report| report.venue_position_id);
4002 let position_ids = match venue_position_id {
4003 Some(open_position_id) => match cached_positions.as_slice() {
4004 [position] => Some((Some(position.id), Some(open_position_id))),
4005 _ => {
4006 log::warn!(
4007 "Deferring hedge cross-zero reconciliation for {instrument_id}/{account_id}: cached and venue position identities are ambiguous"
4008 );
4009 None
4010 }
4011 },
4012 None => Some((None, None)),
4013 };
4014
4015 position_ids.and_then(|(close_position_id, open_position_id)| {
4016 self.reconcile_cross_zero_position(
4017 &instrument,
4018 account_id,
4019 instrument_id,
4020 cached_signed_qty,
4021 cached_avg_px,
4022 venue_signed_qty,
4023 venue_avg_px,
4024 close_position_id,
4025 open_position_id,
4026 ts_now,
4027 venue_ts_last,
4028 )
4029 })
4030 } else {
4031 let qty_diff = venue_signed_qty - cached_signed_qty;
4032 let order_side = if qty_diff > Decimal::ZERO {
4033 OrderSide::Buy
4034 } else {
4035 OrderSide::Sell
4036 };
4037
4038 let reconciliation_px = calculate_reconciliation_price(
4039 cached_signed_qty,
4040 cached_avg_px,
4041 venue_signed_qty,
4042 venue_avg_px,
4043 );
4044
4045 match reconciliation_px.or(venue_avg_px).or(cached_avg_px) {
4046 Some(fill_px) => {
4047 let fill_qty = qty_diff.abs();
4048 let venue_position_id = venue_report
4049 .as_ref()
4050 .and_then(|report| report.venue_position_id);
4051 let venue_ts_last = venue_report.as_ref().map_or(ts_now, |r| r.ts_last);
4052 Quantity::from_decimal_dp(fill_qty, instrument.size_precision())
4053 .ok()
4054 .map(|order_qty| {
4055 let fill_price =
4056 Price::from_decimal_dp(fill_px, instrument.price_precision()).ok();
4057 let venue_order_id = create_position_reconciliation_venue_order_id(
4058 account_id,
4059 instrument_id,
4060 order_side,
4061 OrderType::Market,
4062 order_qty,
4063 fill_price,
4064 venue_position_id,
4065 None,
4066 venue_ts_last,
4067 );
4068
4069 let mut order_report = OrderStatusReport::new(
4070 account_id,
4071 instrument_id,
4072 None,
4073 venue_order_id,
4074 order_side.into(),
4075 OrderType::Market,
4076 TimeInForce::Gtc,
4077 OrderStatus::Filled,
4078 order_qty,
4079 order_qty,
4080 ts_now,
4081 ts_now,
4082 ts_now,
4083 None,
4084 )
4085 .with_avg_px(fill_px);
4086
4087 if let Some(venue_position_id) = venue_position_id {
4088 order_report =
4089 order_report.with_venue_position_id(venue_position_id);
4090 }
4091
4092 order_report
4093 })
4094 .map(|order_report| {
4095 log::info!(
4096 color = LogColor::Blue as u8;
4097 "Generating synthetic fill for position reconciliation {instrument_id}: side={order_side:?}, qty={}, px={fill_px}", qty_diff.abs(),
4098 );
4099
4100 let (events, _) = self.handle_external_order(
4101 &order_report,
4102 account_id,
4103 &instrument,
4104 &[],
4105 true,
4106 None,
4107 None,
4108 );
4109 events
4110 })
4111 }
4112 None => None,
4113 }
4114 };
4115
4116 if result.is_none() || result.as_ref().is_some_and(Vec::is_empty) {
4118 let new_retries = retries + 1;
4119 self.set_position_reconciliation_retries(key, report_shape, new_retries);
4120 if new_retries >= self.config.position_check_retries {
4121 log::error!(
4122 "Position discrepancy for {} unresolved after {} attempts \
4123 (cached_qty={}, venue_qty={}); \
4124 no further reconciliation attempts will be made for the current report shape",
4125 instrument_id,
4126 self.config.position_check_retries,
4127 cached_signed_qty,
4128 venue_signed_qty,
4129 );
4130 }
4131 } else {
4132 self.position_reconciliation_states.shift_remove(&key);
4133 }
4134
4135 result
4136 }
4137
4138 fn set_position_reconciliation_retries(
4139 &mut self,
4140 key: InstrumentAccountKey,
4141 report_shape: PositionReportShape,
4142 retries: u32,
4143 ) {
4144 self.position_reconciliation_states.insert(
4145 key,
4146 PositionReconciliationState {
4147 report_shape,
4148 retries,
4149 },
4150 );
4151 }
4152
4153 fn position_qty_aggregates(
4154 signed_quantities: impl Iterator<Item = Decimal>,
4155 ) -> (Decimal, Decimal, Decimal) {
4156 signed_quantities.fold(
4157 (Decimal::ZERO, Decimal::ZERO, Decimal::ZERO),
4158 |(net, long, short), qty| {
4159 if qty > Decimal::ZERO {
4160 (net + qty, long + qty, short)
4161 } else {
4162 (net + qty, long, short + qty.abs())
4163 }
4164 },
4165 )
4166 }
4167
4168 #[expect(clippy::too_many_arguments)]
4171 fn reconcile_cross_zero_position(
4172 &self,
4173 instrument: &InstrumentAny,
4174 account_id: AccountId,
4175 instrument_id: InstrumentId,
4176 cached_signed_qty: Decimal,
4177 cached_avg_px: Option<Decimal>,
4178 venue_signed_qty: Decimal,
4179 venue_avg_px: Option<Decimal>,
4180 close_position_id: Option<PositionId>,
4181 open_position_id: Option<PositionId>,
4182 ts_now: UnixNanos,
4183 venue_ts_last: UnixNanos,
4184 ) -> Option<Vec<OrderEventAny>> {
4185 log::info!(
4186 color = LogColor::Blue as u8;
4187 "Position crosses zero for {instrument_id}: cached={cached_signed_qty}, venue={venue_signed_qty}. Splitting into two fills",
4188 );
4189
4190 let close_qty = cached_signed_qty.abs();
4191 let close_side = if cached_signed_qty < Decimal::ZERO {
4192 OrderSide::Buy } else {
4194 OrderSide::Sell };
4196 let open_qty = venue_signed_qty.abs();
4197 let open_side = if venue_signed_qty > Decimal::ZERO {
4198 OrderSide::Buy } else {
4200 OrderSide::Sell };
4202
4203 let Some(close_px) = cached_avg_px else {
4204 log::warn!("Cannot close position for {instrument_id}: no cached average price");
4205 return None;
4206 };
4207
4208 let open_report = match venue_avg_px {
4209 Some(open_px) => Some((
4210 build_cross_zero_leg_report(
4211 instrument,
4212 account_id,
4213 instrument_id,
4214 open_side,
4215 open_qty,
4216 open_px,
4217 open_position_id,
4218 "OPEN",
4219 ts_now,
4220 venue_ts_last,
4221 )?,
4222 open_px,
4223 )),
4224 None => None,
4225 };
4226
4227 let close_report = build_cross_zero_leg_report(
4228 instrument,
4229 account_id,
4230 instrument_id,
4231 close_side,
4232 close_qty,
4233 close_px,
4234 close_position_id,
4235 "CLOSE",
4236 ts_now,
4237 venue_ts_last,
4238 )?;
4239
4240 log::info!(
4241 color = LogColor::Blue as u8;
4242 "Generating close fill for cross-zero {instrument_id}: side={close_side:?}, qty={close_qty}, px={close_px}",
4243 );
4244
4245 let (close_events, _) = self.handle_external_order(
4246 &close_report,
4247 account_id,
4248 instrument,
4249 &[],
4250 true,
4251 None,
4252 None,
4253 );
4254 let mut all_events = close_events;
4255
4256 if let Some((open_report, open_px)) = open_report {
4257 log::info!(
4258 color = LogColor::Blue as u8;
4259 "Generating open fill for cross-zero {instrument_id}: side={open_side:?}, qty={open_qty}, px={open_px}",
4260 );
4261
4262 let (open_events, _) = self.handle_external_order(
4263 &open_report,
4264 account_id,
4265 instrument,
4266 &[],
4267 true,
4268 None,
4269 None,
4270 );
4271 all_events.extend(open_events);
4272 } else {
4273 log::warn!("Cannot open new position for {instrument_id}: no venue average price");
4274 }
4275
4276 Some(all_events)
4277 }
4278
4279 fn create_position_from_report(
4284 &self,
4285 report: &PositionStatusReport,
4286 account_id: AccountId,
4287 instrument: &InstrumentAny,
4288 ) -> Option<Vec<OrderEventAny>> {
4289 let instrument_id = report.instrument_id;
4290 let venue_signed_qty = report.signed_decimal_qty;
4291
4292 if venue_signed_qty == Decimal::ZERO {
4293 return None;
4294 }
4295
4296 let order_side = if venue_signed_qty > Decimal::ZERO {
4297 OrderSide::Buy
4298 } else {
4299 OrderSide::Sell
4300 };
4301
4302 let qty_abs = venue_signed_qty.abs();
4303 let venue_avg_px = report.avg_px_open?;
4304
4305 let ts_now = self.clock.borrow().timestamp_ns();
4306 let order_qty = Quantity::from_decimal_dp(qty_abs, instrument.size_precision()).ok()?;
4307 let fill_price = Price::from_decimal_dp(venue_avg_px, instrument.price_precision()).ok();
4308 let venue_order_id = create_position_reconciliation_venue_order_id(
4309 account_id,
4310 instrument_id,
4311 order_side,
4312 OrderType::Market,
4313 order_qty,
4314 fill_price,
4315 report.venue_position_id,
4316 None,
4317 report.ts_last,
4318 );
4319
4320 let mut order_report = OrderStatusReport::new(
4321 account_id,
4322 instrument_id,
4323 None,
4324 venue_order_id,
4325 order_side.into(),
4326 OrderType::Market,
4327 TimeInForce::Gtc,
4328 OrderStatus::Filled,
4329 order_qty,
4330 order_qty,
4331 ts_now,
4332 ts_now,
4333 ts_now,
4334 None,
4335 )
4336 .with_avg_px(venue_avg_px);
4337
4338 if let Some(venue_position_id) = report.venue_position_id {
4340 order_report = order_report.with_venue_position_id(venue_position_id);
4341 }
4342
4343 log::info!(
4344 color = LogColor::Blue as u8;
4345 "Creating position from venue report for {instrument_id}: side={order_side:?}, qty={qty_abs}, avg_px={venue_avg_px}",
4346 );
4347
4348 let (events, _) = self.handle_external_order(
4349 &order_report,
4350 account_id,
4351 instrument,
4352 &[],
4353 true,
4354 None,
4355 None,
4356 );
4357 Some(events)
4358 }
4359
4360 fn reconcile_position_report(
4361 &self,
4362 report: &PositionStatusReport,
4363 account_id: AccountId,
4364 instruments_with_unattributed_fills: &IndexSet<InstrumentId>,
4365 ) -> Option<Vec<OrderEventAny>> {
4366 if report.venue_position_id.is_some() {
4367 self.reconcile_position_report_hedging(
4368 report,
4369 account_id,
4370 instruments_with_unattributed_fills,
4371 )
4372 } else {
4373 self.reconcile_position_report_netting(report, account_id)
4374 }
4375 }
4376
4377 fn reconcile_position_report_hedging(
4378 &self,
4379 report: &PositionStatusReport,
4380 account_id: AccountId,
4381 instruments_with_unattributed_fills: &IndexSet<InstrumentId>,
4382 ) -> Option<Vec<OrderEventAny>> {
4383 let venue_position_id = report.venue_position_id?;
4384
4385 if instruments_with_unattributed_fills.contains(&report.instrument_id) {
4388 log::debug!(
4389 "Skipping hedge position {venue_position_id} reconciliation: unattributed fills in batch"
4390 );
4391 return None;
4392 }
4393
4394 log::debug!(
4395 "Reconciling HEDGE position for {}, venue_position_id={}",
4396 report.instrument_id,
4397 venue_position_id
4398 );
4399
4400 let position = {
4401 let cache = self.cache.borrow();
4402 cache.position_owned(&venue_position_id)
4403 };
4404
4405 match position {
4406 Some(position) => {
4407 let cached_signed_qty = position.signed_decimal_qty();
4408 let venue_signed_qty = report.signed_decimal_qty;
4409
4410 if cached_signed_qty == venue_signed_qty {
4411 log::debug!(
4412 "Hedge position {venue_position_id} matches venue: qty={cached_signed_qty}"
4413 );
4414 return None;
4415 }
4416
4417 if venue_signed_qty == Decimal::ZERO && cached_signed_qty == Decimal::ZERO {
4418 return None;
4419 }
4420
4421 if !self.config.generate_missing_orders {
4422 log::error!(
4423 "Cannot reconcile {} {}: position net qty {} != reported net qty {} \
4424 and `generate_missing_orders` is disabled",
4425 report.instrument_id,
4426 venue_position_id,
4427 cached_signed_qty,
4428 venue_signed_qty
4429 );
4430 return None;
4431 }
4432
4433 self.reconcile_hedge_position_discrepancy(
4434 report,
4435 account_id,
4436 &position,
4437 cached_signed_qty,
4438 )
4439 }
4440 None => {
4441 if report.signed_decimal_qty == Decimal::ZERO {
4442 return None;
4443 }
4444
4445 if !self.config.generate_missing_orders {
4446 log::error!(
4447 "Cannot reconcile position: {venue_position_id} not found and `generate_missing_orders` is disabled"
4448 );
4449 return None;
4450 }
4451
4452 self.reconcile_missing_hedge_position(report, account_id)
4453 }
4454 }
4455 }
4456
4457 fn reconcile_hedge_position_discrepancy(
4458 &self,
4459 report: &PositionStatusReport,
4460 account_id: AccountId,
4461 position: &Position,
4462 cached_signed_qty: Decimal,
4463 ) -> Option<Vec<OrderEventAny>> {
4464 let instrument = self.get_instrument(&report.instrument_id)?;
4465 let venue_signed_qty = report.signed_decimal_qty;
4466
4467 let diff = (cached_signed_qty - venue_signed_qty).abs();
4468 let diff_qty = Quantity::from_decimal_dp(diff, instrument.size_precision()).ok()?;
4469
4470 if diff_qty.is_zero() {
4471 log::debug!(
4472 "Difference quantity rounds to zero for {}, skipping",
4473 instrument.id()
4474 );
4475 return None;
4476 }
4477
4478 let venue_position_id = report.venue_position_id?;
4479 log::warn!(
4480 "Hedge position discrepancy for {} {}: cached={}, venue={}, generating reconciliation order",
4481 report.instrument_id,
4482 venue_position_id,
4483 cached_signed_qty,
4484 venue_signed_qty
4485 );
4486
4487 let current_avg_px = if position.avg_px_open > 0.0 {
4488 Decimal::from_str(&position.avg_px_open.to_string()).ok()
4489 } else {
4490 None
4491 };
4492
4493 self.create_position_reconciliation_order(
4494 report,
4495 account_id,
4496 &instrument,
4497 cached_signed_qty,
4498 diff_qty,
4499 current_avg_px,
4500 )
4501 }
4502
4503 fn reconcile_missing_hedge_position(
4504 &self,
4505 report: &PositionStatusReport,
4506 account_id: AccountId,
4507 ) -> Option<Vec<OrderEventAny>> {
4508 let instrument = self.get_instrument(&report.instrument_id)?;
4509 let venue_signed_qty = report.signed_decimal_qty;
4510
4511 let qty = venue_signed_qty.abs();
4512 let diff_qty = Quantity::from_decimal_dp(qty, instrument.size_precision()).ok()?;
4513
4514 if diff_qty.is_zero() {
4515 return None;
4516 }
4517
4518 let venue_position_id = report.venue_position_id?;
4519 log::warn!(
4520 "Missing hedge position for {} {}: venue reports {}, generating reconciliation order",
4521 report.instrument_id,
4522 venue_position_id,
4523 venue_signed_qty
4524 );
4525
4526 self.create_position_reconciliation_order(
4527 report,
4528 account_id,
4529 &instrument,
4530 Decimal::ZERO,
4531 diff_qty,
4532 None,
4533 )
4534 }
4535
4536 fn reconcile_position_report_netting(
4537 &self,
4538 report: &PositionStatusReport,
4539 account_id: AccountId,
4540 ) -> Option<Vec<OrderEventAny>> {
4541 let instrument_id = report.instrument_id;
4542
4543 log::debug!("Reconciling NET position for {instrument_id}");
4544
4545 let instrument = self.get_instrument(&instrument_id)?;
4546
4547 let (cached_signed_qty, cached_avg_px) = {
4548 let cache = self.cache.borrow();
4549 let positions =
4550 cache.positions_open(None, Some(&instrument_id), None, Some(&account_id), None);
4551
4552 if positions.is_empty() {
4553 (Decimal::ZERO, None)
4554 } else {
4555 let mut total_signed_qty = Decimal::ZERO;
4556 let mut total_value = Decimal::ZERO;
4557 let mut total_qty = Decimal::ZERO;
4558
4559 for pos in positions {
4560 total_signed_qty += pos.signed_decimal_qty();
4561 let qty = pos.signed_decimal_qty().abs();
4562 if pos.avg_px_open > 0.0
4563 && qty > Decimal::ZERO
4564 && let Ok(avg_px) = Decimal::from_str(&pos.avg_px_open.to_string())
4565 {
4566 total_value += avg_px * qty;
4567 total_qty += qty;
4568 }
4569 }
4570
4571 let avg_px = if total_qty > Decimal::ZERO {
4572 Some(total_value / total_qty)
4573 } else {
4574 None
4575 };
4576
4577 (total_signed_qty, avg_px)
4578 }
4579 };
4580
4581 let venue_signed_qty = report.signed_decimal_qty;
4582
4583 log::debug!("venue_signed_qty={venue_signed_qty}, cached_signed_qty={cached_signed_qty}");
4584
4585 let tolerance = self.position_reconciliation_tolerance(account_id);
4586 if (cached_signed_qty - venue_signed_qty).abs() <= tolerance {
4587 log::debug!("Position quantities match for {instrument_id}, no reconciliation needed");
4588 return None;
4589 }
4590
4591 if !self.config.generate_missing_orders {
4592 log::debug!(
4593 "Discrepancy for {instrument_id} position when `generate_missing_orders` disabled, skipping"
4594 );
4595 return None;
4596 }
4597
4598 let diff = (cached_signed_qty - venue_signed_qty).abs();
4599 let diff_qty = Quantity::from_decimal_dp(diff, instrument.size_precision()).ok()?;
4600
4601 if diff_qty.is_zero() {
4602 log::debug!(
4603 "Difference quantity rounds to zero for {instrument_id}, skipping order generation"
4604 );
4605 return None;
4606 }
4607
4608 let crosses_zero = cached_signed_qty != Decimal::ZERO
4609 && venue_signed_qty != Decimal::ZERO
4610 && ((cached_signed_qty > Decimal::ZERO && venue_signed_qty < Decimal::ZERO)
4611 || (cached_signed_qty < Decimal::ZERO && venue_signed_qty > Decimal::ZERO));
4612
4613 if crosses_zero {
4614 let ts_now = self.clock.borrow().timestamp_ns();
4615 return self.reconcile_cross_zero_position(
4616 &instrument,
4617 account_id,
4618 instrument_id,
4619 cached_signed_qty,
4620 cached_avg_px,
4621 venue_signed_qty,
4622 report.avg_px_open,
4623 None,
4624 None,
4625 ts_now,
4626 report.ts_last,
4627 );
4628 }
4629
4630 if cached_signed_qty == Decimal::ZERO {
4631 return self.create_position_from_report(report, account_id, &instrument);
4632 }
4633
4634 self.create_position_reconciliation_order(
4635 report,
4636 account_id,
4637 &instrument,
4638 cached_signed_qty,
4639 diff_qty,
4640 cached_avg_px,
4641 )
4642 }
4643
4644 fn create_position_reconciliation_order(
4645 &self,
4646 report: &PositionStatusReport,
4647 account_id: AccountId,
4648 instrument: &InstrumentAny,
4649 cached_signed_qty: Decimal,
4650 diff_qty: Quantity,
4651 current_avg_px: Option<Decimal>,
4652 ) -> Option<Vec<OrderEventAny>> {
4653 let venue_signed_qty = report.signed_decimal_qty;
4654 let instrument_id = report.instrument_id;
4655
4656 let order_side = if venue_signed_qty > cached_signed_qty {
4657 OrderSide::Buy
4658 } else {
4659 OrderSide::Sell
4660 };
4661
4662 let reconciliation_px = calculate_reconciliation_price(
4663 cached_signed_qty,
4664 current_avg_px,
4665 venue_signed_qty,
4666 report.avg_px_open,
4667 );
4668
4669 let fill_px = reconciliation_px
4670 .or(report.avg_px_open)
4671 .or(current_avg_px)?;
4672
4673 let ts_now = self.clock.borrow().timestamp_ns();
4674 let fill_price = Price::from_decimal_dp(fill_px, instrument.price_precision()).ok();
4675 let venue_order_id = create_position_reconciliation_venue_order_id(
4676 account_id,
4677 instrument_id,
4678 order_side,
4679 OrderType::Market,
4680 diff_qty,
4681 fill_price,
4682 report.venue_position_id,
4683 None,
4684 report.ts_last,
4685 );
4686
4687 let mut order_report = OrderStatusReport::new(
4688 account_id,
4689 instrument_id,
4690 None,
4691 venue_order_id,
4692 order_side.into(),
4693 OrderType::Market,
4694 TimeInForce::Gtc,
4695 OrderStatus::Filled,
4696 diff_qty,
4697 diff_qty,
4698 ts_now,
4699 ts_now,
4700 ts_now,
4701 None,
4702 )
4703 .with_avg_px(fill_px);
4704
4705 if let Some(venue_position_id) = report.venue_position_id {
4706 order_report = order_report.with_venue_position_id(venue_position_id);
4707 }
4708
4709 log::info!(
4710 color = LogColor::Blue as u8;
4711 "Generating reconciliation order for {instrument_id}: side={order_side:?}, qty={diff_qty}, px={fill_px}",
4712 );
4713
4714 let (events, _) = self.handle_external_order(
4715 &order_report,
4716 account_id,
4717 instrument,
4718 &[],
4719 true,
4720 None,
4721 None,
4722 );
4723 Some(events)
4724 }
4725
4726 fn reconcile_order_report(
4727 &self,
4728 order: &OrderAny,
4729 report: &OrderStatusReport,
4730 instrument: Option<&InstrumentAny>,
4731 commission_client: Option<&dyn ExecutionClient>,
4732 ) -> anyhow::Result<Option<OrderEventAny>> {
4733 let ts_now = self.clock.borrow().timestamp_ns();
4734 let commission = if matches!(
4735 report.order_status,
4736 OrderStatus::PartiallyFilled | OrderStatus::Filled
4737 ) && report.filled_qty > order.filled_qty()
4738 {
4739 let Some(instrument) = instrument else {
4740 return Ok(reconcile_order_report_with_commission(
4741 order, report, None, ts_now, None,
4742 ));
4743 };
4744 let fill_qty = report.filled_qty - order.filled_qty();
4745 Self::resolve_inferred_fill_commission(
4746 commission_client,
4747 instrument,
4748 fill_qty,
4749 incremental_inferred_fill_price_and_liquidity(order, report, instrument),
4750 )?
4751 } else {
4752 None
4753 };
4754
4755 Ok(reconcile_order_report_with_commission(
4756 order, report, instrument, ts_now, commission,
4757 ))
4758 }
4759
4760 fn reconcile_order_with_fills(
4762 &mut self,
4763 order: &OrderAny,
4764 report: &OrderStatusReport,
4765 fills: &[&FillReport],
4766 instrument: Option<&InstrumentAny>,
4767 fill_queue: &mut ReconciliationFillQueue,
4768 commission_client: Option<&dyn ExecutionClient>,
4769 ) -> Vec<OrderEventAny> {
4770 let mut events = Vec::new();
4771 let mut working = order.clone();
4772 let mut sorted_fills: Vec<&FillReport> = fills.to_vec();
4773 sorted_fills.sort_by_key(|f| f.ts_event);
4774
4775 let ts_now = self.clock.borrow().timestamp_ns();
4776
4777 if matches!(
4778 report.order_status,
4779 OrderStatus::Canceled | OrderStatus::Expired
4780 ) && report.ts_triggered.is_some()
4781 && working.status() != OrderStatus::Triggered
4782 && TRIGGERABLE_ORDER_TYPES.contains(&working.order_type())
4783 {
4784 let triggered = create_reconciliation_triggered(&working, report, ts_now);
4785 if working.apply(triggered.clone()).is_ok() {
4786 events.push(triggered);
4787 }
4788 }
4789
4790 let requires_snapshot_projection = !sorted_fills.is_empty()
4791 || report.order_status == OrderStatus::Voided
4792 || report.filled_qty < working.filled_qty();
4793 if !requires_snapshot_projection {
4794 match self.reconcile_order_report(&working, report, instrument, commission_client) {
4795 Ok(Some(event)) => events.push(event),
4796 Ok(None) => {}
4797 Err(e) => log::error!(
4798 "Deferring inferred fill for {}: venue commission calculation failed: {e}",
4799 order.client_order_id(),
4800 ),
4801 }
4802 return events;
4803 }
4804
4805 for event in generate_reconciliation_order_pre_fill_events(&working, report, ts_now) {
4806 if let Err(e) = working.apply(event.clone()) {
4807 log::warn!(
4808 "Cannot project reconciliation event for {}: {e}",
4809 order.client_order_id()
4810 );
4811 return events;
4812 }
4813 events.push(event);
4814 }
4815
4816 if let Some(inst) = instrument {
4817 for fill in sorted_fills {
4818 let Some((event, fill_key)) =
4819 self.create_order_fill(&working, fill, inst, &fill_queue.pending_fill_keys)
4820 else {
4821 continue;
4822 };
4823
4824 if let Err(e) = working.apply(event.clone()) {
4825 if let OrderEventAny::Filled(fill) = &event
4826 && self.is_fill_applied(fill, fill_key)
4827 {
4828 self.processed_fills.mark(fill_key);
4829 } else {
4830 log::warn!(
4831 "Cannot project reconciliation fill for {}: {e}",
4832 order.client_order_id()
4833 );
4834 }
4835 return events;
4836 }
4837 fill_queue.push(&mut events, event, fill_key);
4838 }
4839 }
4840
4841 let commission = if report.filled_qty > working.filled_qty()
4842 && let Some(instrument) = instrument
4843 {
4844 let fill_qty = report.filled_qty - working.filled_qty();
4845
4846 match Self::resolve_inferred_fill_commission(
4847 commission_client,
4848 instrument,
4849 fill_qty,
4850 incremental_inferred_fill_price_and_liquidity(&working, report, instrument),
4851 ) {
4852 Ok(commission) => commission,
4853 Err(e) => {
4854 log::error!(
4855 "Deferring inferred fill for {}: venue commission calculation failed: {e}",
4856 order.client_order_id(),
4857 );
4858 return events;
4859 }
4860 }
4861 } else {
4862 None
4863 };
4864
4865 for event in generate_reconciliation_order_snapshot_events_with_commission(
4866 &working, report, instrument, ts_now, commission,
4867 ) {
4868 if let Err(e) = working.apply(event.clone()) {
4869 log::warn!(
4870 "Cannot project reconciliation snapshot event for {}: {e}",
4871 order.client_order_id()
4872 );
4873 break;
4874 }
4875 events.push(event);
4876 }
4877
4878 events
4879 }
4880
4881 fn resolve_inferred_fill_commission(
4882 client: Option<&dyn ExecutionClient>,
4883 instrument: &InstrumentAny,
4884 fill_qty: Quantity,
4885 price_and_liquidity: Option<(Price, LiquiditySide)>,
4886 ) -> anyhow::Result<Option<Money>> {
4887 let Some(client) = client else {
4888 anyhow::bail!("responsible execution client is unavailable");
4889 };
4890 let Some((last_px, liquidity_side)) = price_and_liquidity else {
4891 return Ok(None);
4892 };
4893
4894 client.calculate_commission(instrument, fill_qty, last_px, liquidity_side)
4895 }
4896
4897 #[expect(clippy::too_many_arguments)]
4898 fn handle_external_order(
4899 &self,
4900 report: &OrderStatusReport,
4901 account_id: AccountId,
4902 instrument: &InstrumentAny,
4903 fills: &[&FillReport],
4904 is_synthetic: bool,
4905 fill_queue: Option<&mut ReconciliationFillQueue>,
4906 commission_client: Option<&dyn ExecutionClient>,
4907 ) -> (Vec<OrderEventAny>, Option<ExternalOrderMetadata>) {
4908 let (strategy_id, tags) =
4909 if let Some(claimed_strategy) = self.external_order_claims.get(&report.instrument_id) {
4910 let order_id = report
4911 .client_order_id
4912 .map_or_else(|| report.venue_order_id.to_string(), |id| id.to_string());
4913 log::info!(
4914 color = LogColor::Blue as u8;
4915 "External order {} for {} claimed by strategy {}",
4916 order_id,
4917 report.instrument_id,
4918 claimed_strategy,
4919 );
4920 (*claimed_strategy, None)
4921 } else {
4922 let tag = if is_synthetic {
4924 *TAG_RECONCILIATION
4925 } else {
4926 *TAG_VENUE
4927 };
4928 (StrategyId::from("EXTERNAL"), Some(vec![tag]))
4929 };
4930
4931 if self.config.filter_unclaimed_external && !is_synthetic {
4933 return (Vec::new(), None);
4934 }
4935
4936 let client_order_id = report
4937 .client_order_id
4938 .unwrap_or_else(|| ClientOrderId::from(report.venue_order_id.as_str()));
4939
4940 if !report.quantity.is_positive() {
4941 log::error!(
4942 "Skipping external order {} ({}) for {}: non-positive quantity in report {:?}",
4943 client_order_id,
4944 report.venue_order_id,
4945 report.instrument_id,
4946 report,
4947 );
4948 return (Vec::new(), None);
4949 }
4950
4951 let ts_now = self.clock.borrow().timestamp_ns();
4952 let Some(order_side) = report.order_side else {
4953 log::error!(
4954 "Skipping external order {} ({}) for {}: order side is not specified",
4955 client_order_id,
4956 report.venue_order_id,
4957 report.instrument_id,
4958 );
4959 return (Vec::new(), None);
4960 };
4961
4962 let initialized = match OrderInitialized::new_checked(
4963 self.config.trader_id,
4964 strategy_id,
4965 report.instrument_id,
4966 client_order_id,
4967 order_side,
4968 report.order_type,
4969 report.quantity,
4970 report.time_in_force,
4971 report.post_only,
4972 report.reduce_only,
4973 false, true, UUID4::new(),
4976 ts_now,
4977 ts_now,
4978 report.price,
4979 report.activation_price,
4980 report.trigger_price,
4981 report.trigger_type,
4982 report.limit_offset,
4983 report.trailing_offset,
4984 report.trailing_offset_type,
4985 report.expire_time,
4986 report.display_qty,
4987 None, None, report.contingency_type,
4990 report.order_list_id,
4991 report.linked_order_ids.clone(),
4992 report.parent_order_id,
4993 None, None, None, tags,
4997 ) {
4998 Ok(initialized) => initialized,
4999 Err(e) => {
5000 log::error!("Failed to create order from report: {e}");
5001 return (Vec::new(), None);
5002 }
5003 };
5004
5005 let initialized = OrderEventAny::Initialized(initialized);
5006 let order = match OrderAny::from_events(vec![initialized.clone()]) {
5007 Ok(order) => order,
5008 Err(e) => {
5009 log::error!("Failed to create order from report: {e}");
5010 return (Vec::new(), None);
5011 }
5012 };
5013
5014 let replace_inferred_fill = !fills.is_empty()
5015 && matches!(
5016 report.order_status,
5017 OrderStatus::Canceled
5018 | OrderStatus::Expired
5019 | OrderStatus::Filled
5020 | OrderStatus::PartiallyFilled
5021 );
5022 let mut prepared_fills = Vec::new();
5023 let mut prepared_fill_keys = fill_queue
5024 .as_deref()
5025 .map(|queue| queue.pending_fill_keys.clone())
5026 .unwrap_or_default();
5027 let mut real_fill_total = Decimal::ZERO;
5028
5029 if replace_inferred_fill {
5030 let mut sorted_fills: Vec<&FillReport> = fills.to_vec();
5031 sorted_fills.sort_by_key(|fill| fill.ts_event);
5032 if fill_queue.is_none() {
5033 log::error!(
5034 "Cannot reconcile external order {client_order_id}: fill queue is unavailable"
5035 );
5036 return (Vec::new(), None);
5037 }
5038
5039 for fill in sorted_fills {
5040 if let Some((fill_event, fill_key)) =
5041 self.create_order_fill(&order, fill, instrument, &prepared_fill_keys)
5042 {
5043 real_fill_total += fill.last_qty.as_decimal();
5044 prepared_fill_keys.insert(fill_key);
5045 prepared_fills.push((fill_event, fill_key));
5046 }
5047 }
5048 }
5049
5050 let report_filled = report.filled_qty.as_decimal();
5051 let inferred_qty = if report_filled.is_zero() {
5052 None
5053 } else if replace_inferred_fill {
5054 if real_fill_total < report_filled {
5055 match Quantity::from_decimal_dp(
5056 report_filled - real_fill_total,
5057 instrument.size_precision(),
5058 ) {
5059 Ok(quantity) => Some(quantity),
5060 Err(e) => {
5061 log::error!(
5062 "Cannot reconcile external order {client_order_id}: residual fill quantity is invalid: {e}"
5063 );
5064 return (Vec::new(), None);
5065 }
5066 }
5067 } else {
5068 None
5069 }
5070 } else if matches!(
5071 report.order_status,
5072 OrderStatus::PartiallyFilled
5073 | OrderStatus::Filled
5074 | OrderStatus::Canceled
5075 | OrderStatus::Expired
5076 | OrderStatus::Voided
5077 ) {
5078 Some(report.filled_qty)
5079 } else {
5080 None
5081 };
5082
5083 let inferred_commission = if is_synthetic {
5084 None
5085 } else if let Some(inferred_qty) = inferred_qty {
5086 match Self::resolve_inferred_fill_commission(
5087 commission_client,
5088 instrument,
5089 inferred_qty,
5090 inferred_fill_price_and_liquidity(&order, report, instrument),
5091 ) {
5092 Ok(commission) => commission,
5093 Err(e) => {
5094 log::error!(
5095 "Deferring external order {client_order_id}: venue commission calculation failed: {e}"
5096 );
5097 return (Vec::new(), None);
5098 }
5099 }
5100 } else {
5101 None
5102 };
5103
5104 {
5105 let mut cache = self.cache.borrow_mut();
5106 let source_client_id = if is_synthetic {
5107 None
5108 } else {
5109 commission_client.map(ExecutionClient::client_id)
5110 };
5111
5112 if let Err(e) = cache.add_order(order.clone(), None, source_client_id, false) {
5113 match cache.order(&client_order_id) {
5117 Some(existing) if is_synthetic && existing.is_closed() => {
5118 log::debug!(
5119 "Skipping synthetic reconciliation order {client_order_id} for {}: \
5120 replay deduped (cached status={:?})",
5121 report.instrument_id,
5122 existing.status(),
5123 );
5124 }
5125 Some(existing) if is_synthetic => {
5126 log::warn!(
5127 "Synthetic reconciliation order {client_order_id} for {} exists in \
5128 cache in non-terminal state {:?}; fill not regenerated",
5129 report.instrument_id,
5130 existing.status(),
5131 );
5132 }
5133 _ => {
5134 log::error!("Failed to add external order to cache: {e}");
5135 }
5136 }
5137 return (Vec::new(), None);
5138 }
5139
5140 if let Err(e) = cache.index_venue_order_id(&client_order_id, &report.venue_order_id) {
5141 log::warn!("Failed to index venue order ID: {e}");
5142 }
5143 }
5144
5145 Self::publish_order_event(&initialized);
5146
5147 log::info!(
5148 color = LogColor::Blue as u8;
5149 "Created external order {} ({}) for {} [{}]",
5150 client_order_id,
5151 report.venue_order_id,
5152 report.instrument_id,
5153 report.order_status,
5154 );
5155
5156 let ts_now = self.clock.borrow().timestamp_ns();
5157 let mut order_events = generate_external_order_status_events_with_commission(
5158 &order,
5159 report,
5160 &account_id,
5161 instrument,
5162 ts_now,
5163 inferred_commission,
5164 );
5165
5166 if replace_inferred_fill {
5167 let terminal_event = if order_events.last().is_some_and(|event| {
5168 matches!(
5169 event,
5170 OrderEventAny::Canceled(_) | OrderEventAny::Expired(_),
5171 )
5172 }) {
5173 order_events.pop()
5174 } else {
5175 None
5176 };
5177
5178 if order_events
5179 .last()
5180 .is_some_and(|event| matches!(event, OrderEventAny::Filled(_)))
5181 {
5182 order_events.pop();
5183 }
5184
5185 let fill_queue =
5186 fill_queue.expect("fill queue availability was checked before cache mutation");
5187 for (fill_event, fill_key) in prepared_fills {
5188 fill_queue.push(&mut order_events, fill_event, fill_key);
5189 }
5190
5191 if let Some(inferred_qty) = inferred_qty
5192 && let Some(inferred_fill) = create_inferred_fill_for_qty(
5193 &order,
5194 report,
5195 &account_id,
5196 instrument,
5197 inferred_qty,
5198 ts_now,
5199 inferred_commission,
5200 )
5201 {
5202 order_events.push(inferred_fill);
5203 }
5204
5205 if let Some(event) = terminal_event {
5206 order_events.push(event);
5207 }
5208 }
5209
5210 let metadata = ExternalOrderMetadata {
5211 client_order_id,
5212 venue_order_id: report.venue_order_id,
5213 instrument_id: report.instrument_id,
5214 strategy_id,
5215 ts_init: ts_now,
5216 };
5217
5218 (order_events, Some(metadata))
5219 }
5220
5221 fn publish_order_event(event: &OrderEventAny) {
5222 let topic = switchboard::get_event_order_topic(event.strategy_id());
5223 msgbus::publish_order_event(topic, event);
5224 }
5225
5226 fn adjust_mass_status_fills(
5231 &self,
5232 mass_status: &ExecutionMassStatus,
5233 ) -> (
5234 IndexMap<VenueOrderId, OrderStatusReport>,
5235 IndexMap<VenueOrderId, Vec<FillReport>>,
5236 ) {
5237 let mut final_orders: IndexMap<VenueOrderId, OrderStatusReport> =
5238 mass_status.order_reports();
5239 let mut final_fills: IndexMap<VenueOrderId, Vec<FillReport>> = mass_status.fill_reports();
5240
5241 if mass_status.lookback_start().is_some() {
5242 return (final_orders, final_fills);
5243 }
5244
5245 let mut instruments_to_adjust = Vec::new();
5246
5247 for (instrument_id, position_reports) in mass_status.position_reports() {
5248 if !self.should_reconcile_instrument(&instrument_id) {
5249 log::debug!(
5250 "Skipping fill adjustment for {instrument_id}: not in reconciliation_instrument_ids"
5251 );
5252 continue;
5253 }
5254
5255 let is_hedge_mode = position_reports
5258 .iter()
5259 .any(|r| r.venue_position_id.is_some());
5260
5261 if is_hedge_mode {
5262 log::debug!(
5263 "Skipping fill adjustment for {instrument_id}: hedge mode (has venue_position_id)"
5264 );
5265 continue;
5266 }
5267
5268 let has_retained_position = {
5269 let cache = self.cache.borrow();
5270 !cache
5271 .positions_open(
5272 None,
5273 Some(&instrument_id),
5274 None,
5275 Some(&mass_status.account_id),
5276 None,
5277 )
5278 .is_empty()
5279 };
5280
5281 if has_retained_position {
5282 log::debug!(
5283 "Skipping fill adjustment for {instrument_id}: retained open position in cache"
5284 );
5285 continue;
5286 }
5287
5288 if let Some(instrument) = self.get_instrument(&instrument_id) {
5289 instruments_to_adjust.push(instrument);
5290 } else {
5291 log::debug!(
5292 "Skipping fill adjustment for {instrument_id}: instrument not found in cache"
5293 );
5294 }
5295 }
5296
5297 if instruments_to_adjust.is_empty() {
5298 return (final_orders, final_fills);
5299 }
5300
5301 log_info!(
5302 "Adjusting fills for {} instrument(s) with position reports",
5303 instruments_to_adjust.len(),
5304 color = LogColor::Blue
5305 );
5306
5307 for instrument in &instruments_to_adjust {
5308 let instrument_id = instrument.id();
5309
5310 let result = if self.config.generate_missing_orders {
5311 process_mass_status_for_reconciliation(mass_status, instrument, None)
5312 } else {
5313 process_mass_status_for_reconciliation_without_synthetic_reports(
5314 mass_status,
5315 instrument,
5316 None,
5317 )
5318 };
5319
5320 match result {
5321 Ok(result) => {
5322 final_orders.retain(|_, order| order.instrument_id != instrument_id);
5323 final_fills.retain(|_, fills| {
5324 fills
5325 .first()
5326 .is_none_or(|f| f.instrument_id != instrument_id)
5327 });
5328
5329 for (venue_order_id, order) in result.orders {
5330 final_orders.insert(venue_order_id, order);
5331 }
5332
5333 for (venue_order_id, fills) in result.fills {
5334 final_fills.insert(venue_order_id, fills);
5335 }
5336 }
5337 Err(e) => {
5338 log::warn!("Failed to adjust fills for {instrument_id}: {e}");
5339 }
5340 }
5341 }
5342
5343 log_info!(
5344 "After adjustment: {} order(s), {} fill group(s)",
5345 final_orders.len(),
5346 final_fills.len(),
5347 color = LogColor::Blue
5348 );
5349
5350 (final_orders, final_fills)
5351 }
5352
5353 fn deduplicate_order_reports<'a>(
5358 reports: impl Iterator<Item = &'a OrderStatusReport>,
5359 ) -> IndexMap<VenueOrderId, &'a OrderStatusReport> {
5360 let mut best_reports: IndexMap<VenueOrderId, &'a OrderStatusReport> = IndexMap::new();
5361
5362 for report in reports {
5363 let dominated = best_reports
5364 .get(&report.venue_order_id)
5365 .is_some_and(|existing| Self::is_more_advanced(existing, report));
5366
5367 if !dominated {
5368 best_reports.insert(report.venue_order_id, report);
5369 }
5370 }
5371
5372 best_reports
5373 }
5374
5375 fn is_more_advanced(a: &OrderStatusReport, b: &OrderStatusReport) -> bool {
5376 if a.filled_qty > b.filled_qty {
5377 return true;
5378 }
5379
5380 if a.filled_qty < b.filled_qty {
5381 return false;
5382 }
5383
5384 Self::status_priority(a.order_status) > Self::status_priority(b.order_status)
5386 }
5387
5388 const fn status_priority(status: OrderStatus) -> u8 {
5389 match status {
5390 OrderStatus::Initialized | OrderStatus::Submitted | OrderStatus::Emulated => 0,
5391 OrderStatus::Released | OrderStatus::Denied => 1,
5392 OrderStatus::Accepted | OrderStatus::PendingUpdate | OrderStatus::PendingCancel => 2,
5393 OrderStatus::Triggered => 3,
5394 OrderStatus::PartiallyFilled => 4,
5395 OrderStatus::Canceled | OrderStatus::Expired | OrderStatus::Rejected => 5,
5396 OrderStatus::Filled | OrderStatus::Voided => 6,
5397 }
5398 }
5399
5400 fn is_exact_order_match(order: &OrderAny, report: &OrderStatusReport) -> bool {
5401 order.status() == report.order_status
5402 && order.filled_qty() == report.filled_qty
5403 && !should_reconciliation_update(order, report)
5404 }
5405
5406 fn is_fill_applied(&self, fill: &OrderFilled, fill_key: FillKey) -> bool {
5407 self.get_order(fill.client_order_id)
5408 .or_else(|| self.get_order_by_venue_order_id(fill.venue_order_id))
5409 .is_some_and(|order| {
5410 order.account_id() == Some(fill_key.0)
5411 && order.instrument_id() == fill_key.1
5412 && order.trade_ids().contains(&&fill_key.2)
5413 })
5414 }
5415
5416 fn create_order_fill(
5417 &self,
5418 order: &OrderAny,
5419 fill: &FillReport,
5420 instrument: &InstrumentAny,
5421 pending_fill_keys: &IndexSet<FillKey>,
5422 ) -> Option<(OrderEventAny, FillKey)> {
5423 let fill_key = (fill.account_id, fill.instrument_id, fill.trade_id);
5424 if self.processed_fills.contains_key(&fill_key) || pending_fill_keys.contains(&fill_key) {
5425 return None;
5426 }
5427
5428 let order_side = order.order_side();
5429 if fill.order_side != order_side {
5430 log::warn!(
5431 "Fill side mismatch for {}: cached={:?}, venue={:?}",
5432 order.client_order_id(),
5433 order_side,
5434 fill.order_side,
5435 );
5436 }
5437
5438 let event = OrderEventAny::Filled(OrderFilled::new(
5439 order.trader_id(),
5440 order.strategy_id(),
5441 order.instrument_id(),
5442 order.client_order_id(),
5443 fill.venue_order_id,
5444 fill.account_id,
5445 fill.trade_id,
5446 fill.order_side,
5447 order.order_type(),
5448 fill.last_qty,
5449 fill.last_px,
5450 instrument.quote_currency(),
5451 fill.liquidity_side,
5452 fill.report_id,
5453 fill.ts_event,
5454 self.clock.borrow().timestamp_ns(),
5455 false,
5456 fill.venue_position_id,
5457 Some(fill.commission),
5458 None,
5459 ));
5460
5461 Some((event, fill_key))
5462 }
5463}
5464
5465pub(crate) async fn request_targeted_order_reports(
5466 clients: &[&dyn ExecutionClient],
5467 queries: Vec<TargetedOrderQuery>,
5468 query_delay: Duration,
5469) -> Vec<TargetedOrderReportResult> {
5470 let mut results = Vec::with_capacity(queries.len());
5471 let mut request_count = 0usize;
5472
5473 for query in queries {
5474 let mut report = None;
5475 let mut report_client_id = None;
5476 let mut coverage_complete = true;
5477
5478 for client_id in &query.responsible_clients {
5479 let client_id = *client_id;
5480 let Some(client) = clients
5481 .iter()
5482 .find(|client| client.client_id() == client_id)
5483 else {
5484 coverage_complete = false;
5485 log::warn!(
5486 "Cannot run targeted order status query for {}: execution client {client_id} is unavailable",
5487 query.client_order_id,
5488 );
5489 continue;
5490 };
5491
5492 if request_count > 0 && !query_delay.is_zero() {
5493 dst::time::sleep(query_delay).await;
5494 }
5495 request_count += 1;
5496
5497 match client.generate_order_status_report(&query.command).await {
5498 Ok(Some(candidate)) if targeted_report_matches(&query, &candidate) => {
5499 report = Some(candidate);
5500 report_client_id = Some(client_id);
5501 break;
5502 }
5503 Ok(Some(candidate)) => {
5504 coverage_complete = false;
5505 log::warn!(
5506 "Ignoring mismatched targeted order status report from {client_id} for {}: client_order_id={:?}, venue_order_id={}, instrument_id={}",
5507 query.client_order_id,
5508 candidate.client_order_id,
5509 candidate.venue_order_id,
5510 candidate.instrument_id,
5511 );
5512 }
5513 Ok(None) => {}
5514 Err(e) => {
5515 coverage_complete = false;
5516 log::warn!(
5517 "Failed targeted order status query from {client_id} for {}: {e}",
5518 query.client_order_id,
5519 );
5520 }
5521 }
5522 }
5523
5524 results.push(TargetedOrderReportResult {
5525 client_order_id: query.client_order_id,
5526 client_id: report_client_id,
5527 report,
5528 coverage_complete,
5529 });
5530 }
5531
5532 results
5533}
5534
5535fn targeted_report_matches(query: &TargetedOrderQuery, report: &OrderStatusReport) -> bool {
5536 let instrument_matches = query
5537 .command
5538 .instrument_id
5539 .is_none_or(|instrument_id| report.instrument_id == instrument_id);
5540 let order_matches = report.client_order_id == Some(query.client_order_id)
5541 || query
5542 .command
5543 .venue_order_id
5544 .is_some_and(|venue_order_id| report.venue_order_id == venue_order_id);
5545
5546 instrument_matches && order_matches
5547}
5548
5549#[cfg(test)]
5550mod tests {
5551 use nautilus_common::clock::TestClock;
5552 use nautilus_core::{Params, datetime::NANOSECONDS_IN_SECOND};
5553 use nautilus_execution::reconciliation::generate_reconciliation_order_events;
5554 use nautilus_model::{
5555 accounts::AccountAny,
5556 enums::{LiquiditySide, OmsType, PositionSide},
5557 events::order::spec::{OrderPendingCancelSpec, OrderPendingUpdateSpec, OrderUpdatedSpec},
5558 identifiers::{Symbol, Venue},
5559 instruments::{
5560 CurrencyPair, Instrument,
5561 stubs::{crypto_perpetual_ethusdt, xbtusd_bitmex},
5562 },
5563 orders::{OrderTestBuilder, stubs::TestOrderEventStubs},
5564 types::{AccountBalance, Currency, MarginBalance, Money},
5565 };
5566 use rstest::rstest;
5567 use rust_decimal_macros::dec;
5568
5569 use super::*;
5570 #[cfg(feature = "node")]
5571 use crate::execution::client::LiveExecutionClient;
5572
5573 #[rstest]
5574 fn test_new_validates_open_check_lookback_mins_boundaries() {
5575 let create_manager = |mins| {
5576 ExecutionManager::new(
5577 Rc::new(RefCell::new(TestClock::new())),
5578 Rc::new(RefCell::new(Cache::default())),
5579 ExecutionManagerConfig {
5580 open_check_lookback_mins: Some(mins),
5581 ..Default::default()
5582 },
5583 )
5584 };
5585
5586 assert!(create_manager(307_445_734).is_ok());
5587 assert!(matches!(
5588 create_manager(307_445_735),
5589 Err(ConfigError::Range { field, .. })
5590 if field == "ExecutionManagerConfig.open_check_lookback_mins"
5591 ));
5592 }
5593
5594 #[rstest]
5595 fn test_new_validates_reconciliation_lookback_mins_boundaries() {
5596 let create_manager = |mins| {
5597 ExecutionManager::new(
5598 Rc::new(RefCell::new(TestClock::new())),
5599 Rc::new(RefCell::new(Cache::default())),
5600 ExecutionManagerConfig {
5601 lookback_mins: Some(mins),
5602 ..Default::default()
5603 },
5604 )
5605 };
5606
5607 assert!(create_manager(307_445_734_561_825_860).is_ok());
5608 assert!(matches!(
5609 create_manager(307_445_734_561_825_861),
5610 Err(ConfigError::Range { field, .. })
5611 if field == "ExecutionManagerConfig.lookback_mins"
5612 ));
5613 }
5614
5615 #[rstest]
5616 fn test_new_reports_every_invalid_lookback_field() {
5617 let error = ExecutionManager::new(
5618 Rc::new(RefCell::new(TestClock::new())),
5619 Rc::new(RefCell::new(Cache::default())),
5620 ExecutionManagerConfig {
5621 lookback_mins: Some(307_445_734_561_825_861),
5622 open_check_lookback_mins: Some(307_445_735),
5623 position_check_lookback_mins: 307_445_735,
5624 ..Default::default()
5625 },
5626 )
5627 .expect_err("all lookback fields are out of range");
5628
5629 let ConfigError::Multiple { errors } = error else {
5630 panic!("expected a `Multiple` error, was {error:?}");
5631 };
5632
5633 let fields = errors
5634 .iter()
5635 .map(|e| match e {
5636 ConfigError::Range { field, .. } => field.as_str(),
5637 other => panic!("expected a `Range` error, was {other:?}"),
5638 })
5639 .collect::<Vec<_>>();
5640
5641 assert_eq!(
5642 fields,
5643 [
5644 "ExecutionManagerConfig.lookback_mins",
5645 "ExecutionManagerConfig.open_check_lookback_mins",
5646 "ExecutionManagerConfig.position_check_lookback_mins",
5647 ]
5648 );
5649 }
5650
5651 #[derive(Clone)]
5652 enum CommissionOutcome {
5653 Value(Money),
5654 NoOverride,
5655 Failure,
5656 }
5657
5658 struct CommissionStubClient {
5659 outcome: CommissionOutcome,
5660 seen: RefCell<Option<(Quantity, Price, LiquiditySide)>>,
5661 }
5662
5663 impl CommissionStubClient {
5664 fn new(outcome: CommissionOutcome) -> Self {
5665 Self {
5666 outcome,
5667 seen: RefCell::new(None),
5668 }
5669 }
5670 }
5671
5672 #[async_trait::async_trait(?Send)]
5673 impl ExecutionClient for CommissionStubClient {
5674 fn is_connected(&self) -> bool {
5675 true
5676 }
5677
5678 fn client_id(&self) -> ClientId {
5679 ClientId::from("STUB")
5680 }
5681
5682 fn account_id(&self) -> AccountId {
5683 AccountId::from("STUB-001")
5684 }
5685
5686 fn venue(&self) -> Venue {
5687 Venue::from("STUB")
5688 }
5689
5690 fn oms_type(&self) -> OmsType {
5691 OmsType::Netting
5692 }
5693
5694 fn get_account(&self) -> Option<AccountAny> {
5695 None
5696 }
5697
5698 fn generate_account_state(
5699 &self,
5700 _balances: Vec<AccountBalance>,
5701 _margins: Vec<MarginBalance>,
5702 _reported: bool,
5703 _ts_event: UnixNanos,
5704 _info: Option<Params>,
5705 ) -> anyhow::Result<()> {
5706 Ok(())
5707 }
5708
5709 fn start(&mut self) -> anyhow::Result<()> {
5710 Ok(())
5711 }
5712
5713 fn stop(&mut self) -> anyhow::Result<()> {
5714 Ok(())
5715 }
5716
5717 fn calculate_commission(
5718 &self,
5719 _instrument: &InstrumentAny,
5720 last_qty: Quantity,
5721 last_px: Price,
5722 liquidity_side: LiquiditySide,
5723 ) -> anyhow::Result<Option<Money>> {
5724 *self.seen.borrow_mut() = Some((last_qty, last_px, liquidity_side));
5725
5726 match &self.outcome {
5727 CommissionOutcome::Value(money) => Ok(Some(*money)),
5728 CommissionOutcome::NoOverride => Ok(None),
5729 CommissionOutcome::Failure => {
5730 anyhow::bail!("commission is not representable as Money")
5731 }
5732 }
5733 }
5734 }
5735
5736 struct PositionCoverageStubClient;
5737
5738 #[async_trait::async_trait(?Send)]
5739 impl ExecutionClient for PositionCoverageStubClient {
5740 fn is_connected(&self) -> bool {
5741 true
5742 }
5743
5744 fn client_id(&self) -> ClientId {
5745 ClientId::from("BYBIT")
5746 }
5747
5748 fn account_id(&self) -> AccountId {
5749 AccountId::from("TEST-001")
5750 }
5751
5752 fn venue(&self) -> Venue {
5753 Venue::from("BYBIT")
5754 }
5755
5756 fn oms_type(&self) -> OmsType {
5757 OmsType::Netting
5758 }
5759
5760 fn get_account(&self) -> Option<AccountAny> {
5761 None
5762 }
5763
5764 fn provides_bulk_position_coverage(&self, instrument_id: InstrumentId) -> bool {
5765 !instrument_id.symbol.as_str().ends_with("-SPOT")
5766 }
5767
5768 fn generate_account_state(
5769 &self,
5770 _balances: Vec<AccountBalance>,
5771 _margins: Vec<MarginBalance>,
5772 _reported: bool,
5773 _ts_event: UnixNanos,
5774 _info: Option<Params>,
5775 ) -> anyhow::Result<()> {
5776 Ok(())
5777 }
5778
5779 fn start(&mut self) -> anyhow::Result<()> {
5780 Ok(())
5781 }
5782
5783 fn stop(&mut self) -> anyhow::Result<()> {
5784 Ok(())
5785 }
5786 }
5787
5788 fn commission_fixtures() -> (OrderAny, OrderStatusReport, InstrumentAny) {
5789 let instrument = crypto_perpetual_ethusdt();
5790 let order = OrderTestBuilder::new(OrderType::Limit)
5791 .instrument_id(instrument.id())
5792 .side(OrderSide::Buy)
5793 .quantity(Quantity::from("10.0"))
5794 .price(Price::from("100.00"))
5795 .build();
5796 let report = OrderStatusReport::new(
5797 AccountId::from("STUB-001"),
5798 instrument.id(),
5799 Some(order.client_order_id()),
5800 VenueOrderId::from("V-1"),
5801 OrderSide::Buy.into(),
5802 OrderType::Limit,
5803 TimeInForce::Gtc,
5804 OrderStatus::Filled,
5805 Quantity::from("10.0"),
5806 Quantity::from("10.0"),
5807 UnixNanos::from(1),
5808 UnixNanos::from(1),
5809 UnixNanos::from(1),
5810 None,
5811 )
5812 .with_avg_px(dec!(100.0));
5813
5814 (order, report, InstrumentAny::CryptoPerpetual(instrument))
5815 }
5816
5817 fn cached_commission_fixtures() -> (
5818 ExecutionManager,
5819 Rc<RefCell<Cache>>,
5820 OrderAny,
5821 OrderStatusReport,
5822 InstrumentAny,
5823 ) {
5824 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
5825 let clock = Rc::new(RefCell::new(TestClock::new()));
5826 let cache = Rc::new(RefCell::new(Cache::default()));
5827 cache
5828 .borrow_mut()
5829 .add_instrument(instrument.clone())
5830 .expect("instrument is cacheable");
5831 let client_order_id = ClientOrderId::from("O-COMMISSION-CACHED");
5832 let venue_order_id = VenueOrderId::from("V-COMMISSION-CACHED");
5833 insert_accepted_limit_order(
5834 &cache,
5835 client_order_id,
5836 venue_order_id,
5837 instrument.id(),
5838 ClientId::from("STUB"),
5839 );
5840 let order = cache
5841 .borrow()
5842 .order_owned(&client_order_id)
5843 .expect("accepted order is cached");
5844 let report = OrderStatusReport::new(
5845 AccountId::from("TEST-001"),
5846 instrument.id(),
5847 Some(client_order_id),
5848 venue_order_id,
5849 OrderSide::Buy.into(),
5850 OrderType::Limit,
5851 TimeInForce::Gtc,
5852 OrderStatus::Filled,
5853 Quantity::from("10.0"),
5854 Quantity::from("10.0"),
5855 UnixNanos::from(1),
5856 UnixNanos::from(1),
5857 UnixNanos::from(1),
5858 None,
5859 )
5860 .with_avg_px(dec!(100.0));
5861 let manager =
5862 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default())
5863 .expect("valid config");
5864
5865 (manager, cache, order, report, instrument)
5866 }
5867
5868 #[rstest]
5869 fn test_resolve_inferred_fill_commission_without_client_fails_closed() {
5870 let (order, report, instrument) = commission_fixtures();
5871
5872 let error = ExecutionManager::resolve_inferred_fill_commission(
5873 None,
5874 &instrument,
5875 Quantity::from("5.0"),
5876 inferred_fill_price_and_liquidity(&order, &report, &instrument),
5877 )
5878 .expect_err("a missing responsible client must defer the fill");
5879
5880 assert_eq!(
5881 error.to_string(),
5882 "responsible execution client is unavailable"
5883 );
5884 }
5885
5886 #[rstest]
5887 fn test_resolve_inferred_fill_commission_without_price_uses_generic_path() {
5888 let instrument = crypto_perpetual_ethusdt();
5889 let order = OrderTestBuilder::new(OrderType::Market)
5890 .instrument_id(instrument.id())
5891 .side(OrderSide::Buy)
5892 .quantity(Quantity::from("10.0"))
5893 .build();
5894 let report = OrderStatusReport::new(
5895 AccountId::from("STUB-001"),
5896 instrument.id(),
5897 Some(order.client_order_id()),
5898 VenueOrderId::from("V-1"),
5899 OrderSide::Buy.into(),
5900 OrderType::Market,
5901 TimeInForce::Gtc,
5902 OrderStatus::Filled,
5903 Quantity::from("10.0"),
5904 Quantity::from("10.0"),
5905 UnixNanos::from(1),
5906 UnixNanos::from(1),
5907 UnixNanos::from(1),
5908 None,
5909 );
5910 let client =
5911 CommissionStubClient::new(CommissionOutcome::Value(Money::new(1.0, Currency::USDT())));
5912 let instrument = InstrumentAny::CryptoPerpetual(instrument);
5913
5914 let commission = ExecutionManager::resolve_inferred_fill_commission(
5915 Some(&client),
5916 &instrument,
5917 Quantity::from("5.0"),
5918 inferred_fill_price_and_liquidity(&order, &report, &instrument),
5919 )
5920 .expect("an unresolvable price is not a failure");
5921
5922 assert_eq!(commission, None, "no price means no venue commission");
5923 }
5924
5925 #[rstest]
5926 fn test_resolve_inferred_fill_commission_returns_venue_value() {
5927 let (order, report, instrument) = commission_fixtures();
5928 let expected = Money::new(2.5, Currency::USDT());
5929 let client = CommissionStubClient::new(CommissionOutcome::Value(expected));
5930
5931 let commission = ExecutionManager::resolve_inferred_fill_commission(
5932 Some(&client),
5933 &instrument,
5934 Quantity::from("5.0"),
5935 inferred_fill_price_and_liquidity(&order, &report, &instrument),
5936 )
5937 .expect("a representable commission succeeds");
5938
5939 assert_eq!(commission, Some(expected));
5940 assert_eq!(
5941 *client.seen.borrow(),
5942 Some((
5943 Quantity::from("5.0"),
5944 Price::from("100.00"),
5945 LiquiditySide::NoLiquiditySide,
5946 )),
5947 "the resolver passes the inferred fill quantity, resolved price, and liquidity side"
5948 );
5949 }
5950
5951 #[rstest]
5952 fn test_resolve_inferred_fill_commission_honors_no_override() {
5953 let (order, report, instrument) = commission_fixtures();
5954 let client = CommissionStubClient::new(CommissionOutcome::NoOverride);
5955
5956 let commission = ExecutionManager::resolve_inferred_fill_commission(
5957 Some(&client),
5958 &instrument,
5959 Quantity::from("5.0"),
5960 inferred_fill_price_and_liquidity(&order, &report, &instrument),
5961 )
5962 .expect("no override is not a failure");
5963
5964 assert_eq!(commission, None);
5965 }
5966
5967 #[rstest]
5968 fn test_resolve_inferred_fill_commission_propagates_failure() {
5969 let (order, report, instrument) = commission_fixtures();
5970 let client = CommissionStubClient::new(CommissionOutcome::Failure);
5971
5972 let result = ExecutionManager::resolve_inferred_fill_commission(
5973 Some(&client),
5974 &instrument,
5975 Quantity::from("5.0"),
5976 inferred_fill_price_and_liquidity(&order, &report, &instrument),
5977 );
5978
5979 assert!(result.is_err(), "a venue failure must not become Ok(None)");
5980 }
5981
5982 fn external_report_with_partial_fill(
5983 instrument: &InstrumentAny,
5984 ) -> (OrderStatusReport, FillReport) {
5985 let account_id = AccountId::from("STUB-001");
5986 let venue_order_id = VenueOrderId::from("V-EXT-1");
5987 let report = OrderStatusReport::new(
5988 account_id,
5989 instrument.id(),
5990 None,
5991 venue_order_id,
5992 OrderSide::Buy.into(),
5993 OrderType::Limit,
5994 TimeInForce::Gtc,
5995 OrderStatus::Filled,
5996 Quantity::from("10.0"),
5997 Quantity::from("10.0"),
5998 UnixNanos::from(1),
5999 UnixNanos::from(1),
6000 UnixNanos::from(1),
6001 None,
6002 )
6003 .with_price(Price::from("100.00"))
6004 .with_avg_px(dec!(100.0));
6005
6006 let fill = FillReport::new(
6007 account_id,
6008 instrument.id(),
6009 venue_order_id,
6010 TradeId::from("T-EXT-1"),
6011 OrderSide::Buy,
6012 Quantity::from("4.0"),
6013 Price::from("100.00"),
6014 Money::new(0.1, Currency::USDT()),
6015 LiquiditySide::Taker,
6016 None,
6017 None,
6018 UnixNanos::from(1),
6019 UnixNanos::from(1),
6020 None,
6021 );
6022
6023 (report, fill)
6024 }
6025
6026 fn inferred_fills(events: &[OrderEventAny]) -> Vec<OrderFilled> {
6027 events
6028 .iter()
6029 .filter_map(|event| match event {
6030 OrderEventAny::Filled(filled) if filled.last_qty == Quantity::from("6.0") => {
6031 Some(filled.clone())
6032 }
6033 _ => None,
6034 })
6035 .collect()
6036 }
6037
6038 #[rstest]
6039 fn test_create_orphan_fill_order_report_rejects_mixed_position_ids() {
6040 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
6041 let account_id = AccountId::from("STUB-001");
6042 let venue_order_id = VenueOrderId::from("V-ORPHAN-1");
6043 let first = FillReport::new(
6044 account_id,
6045 instrument.id(),
6046 venue_order_id,
6047 TradeId::from("T-ORPHAN-1"),
6048 OrderSide::Buy,
6049 Quantity::from("1.0"),
6050 Price::from("100.00"),
6051 Money::from("0.10 USDT"),
6052 LiquiditySide::Taker,
6053 None,
6054 Some(PositionId::from("P-LONG")),
6055 UnixNanos::from(1),
6056 UnixNanos::from(1),
6057 None,
6058 );
6059
6060 let second = FillReport::new(
6061 account_id,
6062 instrument.id(),
6063 venue_order_id,
6064 TradeId::from("T-ORPHAN-2"),
6065 OrderSide::Buy,
6066 Quantity::from("1.0"),
6067 Price::from("101.00"),
6068 Money::from("0.10 USDT"),
6069 LiquiditySide::Maker,
6070 None,
6071 Some(PositionId::from("P-SHORT")),
6072 UnixNanos::from(2),
6073 UnixNanos::from(2),
6074 None,
6075 );
6076
6077 let error =
6078 ExecutionManager::create_orphan_fill_order_report(&[&first, &second], &instrument)
6079 .unwrap_err();
6080
6081 assert_eq!(
6082 error.to_string(),
6083 "venue position ID differs across fill group"
6084 );
6085 }
6086
6087 #[rstest]
6088 fn test_handle_external_order_applies_venue_commission_to_inferred_fill() {
6089 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
6090 let clock = Rc::new(RefCell::new(TestClock::new()));
6091 let cache = Rc::new(RefCell::new(Cache::default()));
6092 cache
6093 .borrow_mut()
6094 .add_instrument(instrument.clone())
6095 .expect("instrument is cacheable");
6096 let manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default())
6097 .expect("valid config");
6098 let (report, fill) = external_report_with_partial_fill(&instrument);
6099 let expected = Money::new(2.5, Currency::USDT());
6100 let client = CommissionStubClient::new(CommissionOutcome::Value(expected));
6101 let mut fill_queue = ReconciliationFillQueue::default();
6102
6103 let (events, _) = manager.handle_external_order(
6104 &report,
6105 AccountId::from("STUB-001"),
6106 &instrument,
6107 &[&fill],
6108 false,
6109 Some(&mut fill_queue),
6110 Some(&client),
6111 );
6112
6113 let inferred = inferred_fills(&events);
6114 assert_eq!(inferred.len(), 1, "one inferred fill covers the 6.0 gap");
6115 assert_eq!(inferred[0].commission, Some(expected));
6116 }
6117
6118 #[rstest]
6119 fn test_handle_external_order_skips_inferred_fill_when_commission_fails() {
6120 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
6121 let clock = Rc::new(RefCell::new(TestClock::new()));
6122 let cache = Rc::new(RefCell::new(Cache::default()));
6123 cache
6124 .borrow_mut()
6125 .add_instrument(instrument.clone())
6126 .expect("instrument is cacheable");
6127 let manager =
6128 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default())
6129 .expect("valid config");
6130 let (report, fill) = external_report_with_partial_fill(&instrument);
6131 let client = CommissionStubClient::new(CommissionOutcome::Failure);
6132 let mut fill_queue = ReconciliationFillQueue::default();
6133
6134 let (events, metadata) = manager.handle_external_order(
6135 &report,
6136 AccountId::from("STUB-001"),
6137 &instrument,
6138 &[&fill],
6139 false,
6140 Some(&mut fill_queue),
6141 Some(&client),
6142 );
6143
6144 assert!(events.is_empty());
6145 assert!(metadata.is_none());
6146 assert!(fill_queue.pending_fill_keys.is_empty());
6147 assert!(
6148 cache
6149 .borrow()
6150 .order(&ClientOrderId::from(report.venue_order_id.as_str()))
6151 .is_none(),
6152 "commission failure must precede external order cache mutation"
6153 );
6154
6155 let expected = Money::new(2.5, Currency::USDT());
6156 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
6157 let (retry_events, retry_metadata) = manager.handle_external_order(
6158 &report,
6159 AccountId::from("STUB-001"),
6160 &instrument,
6161 &[&fill],
6162 false,
6163 Some(&mut fill_queue),
6164 Some(&retry_client),
6165 );
6166 let inferred = inferred_fills(&retry_events);
6167
6168 assert!(retry_metadata.is_some());
6169 assert_eq!(inferred.len(), 1);
6170 assert_eq!(inferred[0].commission, Some(expected));
6171 assert_eq!(fill_queue.pending_fill_keys.len(), 1);
6172 assert!(
6173 cache
6174 .borrow()
6175 .order(&ClientOrderId::from(report.venue_order_id.as_str()))
6176 .is_some()
6177 );
6178 }
6179
6180 #[rstest]
6181 fn test_handle_external_order_without_explicit_fills_resolves_commission_before_cache() {
6182 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
6183 let clock = Rc::new(RefCell::new(TestClock::new()));
6184 let cache = Rc::new(RefCell::new(Cache::default()));
6185 cache
6186 .borrow_mut()
6187 .add_instrument(instrument.clone())
6188 .expect("instrument is cacheable");
6189 let manager =
6190 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default())
6191 .expect("valid config");
6192 let (report, _) = external_report_with_partial_fill(&instrument);
6193 let failing_client = CommissionStubClient::new(CommissionOutcome::Failure);
6194
6195 let (failed_events, failed_metadata) = manager.handle_external_order(
6196 &report,
6197 AccountId::from("STUB-001"),
6198 &instrument,
6199 &[],
6200 false,
6201 None,
6202 Some(&failing_client),
6203 );
6204
6205 assert!(failed_events.is_empty());
6206 assert!(failed_metadata.is_none());
6207 assert!(
6208 cache
6209 .borrow()
6210 .order(&ClientOrderId::from(report.venue_order_id.as_str()))
6211 .is_none()
6212 );
6213
6214 let expected = Money::new(4.0, Currency::USDT());
6215 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
6216 let (retry_events, retry_metadata) = manager.handle_external_order(
6217 &report,
6218 AccountId::from("STUB-001"),
6219 &instrument,
6220 &[],
6221 false,
6222 None,
6223 Some(&retry_client),
6224 );
6225 let fills: Vec<_> = retry_events
6226 .iter()
6227 .filter_map(|event| match event {
6228 OrderEventAny::Filled(fill) => Some(fill),
6229 _ => None,
6230 })
6231 .collect();
6232
6233 assert!(retry_metadata.is_some());
6234 assert_eq!(fills.len(), 1);
6235 assert_eq!(fills[0].last_qty, Quantity::from("10.0"));
6236 assert_eq!(fills[0].commission, Some(expected));
6237 }
6238
6239 #[rstest]
6240 fn test_handle_external_order_with_no_override_emits_fill_without_commission() {
6241 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
6242 let clock = Rc::new(RefCell::new(TestClock::new()));
6243 let cache = Rc::new(RefCell::new(Cache::default()));
6244 cache
6245 .borrow_mut()
6246 .add_instrument(instrument.clone())
6247 .expect("instrument is cacheable");
6248 let manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default())
6249 .expect("valid config");
6250 let (report, fill) = external_report_with_partial_fill(&instrument);
6251 let client = CommissionStubClient::new(CommissionOutcome::NoOverride);
6252 let mut fill_queue = ReconciliationFillQueue::default();
6253
6254 let (events, _) = manager.handle_external_order(
6255 &report,
6256 AccountId::from("STUB-001"),
6257 &instrument,
6258 &[&fill],
6259 false,
6260 Some(&mut fill_queue),
6261 Some(&client),
6262 );
6263
6264 let inferred = inferred_fills(&events);
6265 assert_eq!(inferred.len(), 1);
6266 assert_eq!(inferred[0].commission, None);
6267 }
6268
6269 #[rstest]
6270 fn test_cached_reconciliation_applies_explicit_fill_and_defers_failed_residual() {
6271 let (mut manager, _cache, order, mut report, instrument) = cached_commission_fixtures();
6272 report.avg_px = Some(dec!(60.0));
6273 let explicit_fill = FillReport::new(
6274 report.account_id,
6275 report.instrument_id,
6276 report.venue_order_id,
6277 TradeId::from("T-COMMISSION-EXPLICIT"),
6278 OrderSide::Buy,
6279 Quantity::from("4.0"),
6280 Price::from("100.0"),
6281 Money::new(0.25, Currency::USDT()),
6282 LiquiditySide::Taker,
6283 report.client_order_id,
6284 None,
6285 UnixNanos::from(1),
6286 UnixNanos::from(1),
6287 None,
6288 );
6289 let failing_client = CommissionStubClient::new(CommissionOutcome::Failure);
6290 let mut fill_queue = ReconciliationFillQueue::default();
6291
6292 let first_events = manager.reconcile_order_with_fills(
6293 &order,
6294 &report,
6295 &[&explicit_fill],
6296 Some(&instrument),
6297 &mut fill_queue,
6298 Some(&failing_client),
6299 );
6300 let mut working = order;
6301 for event in &first_events {
6302 working
6303 .apply(event.clone())
6304 .expect("explicit fill projects cleanly");
6305 }
6306
6307 assert_eq!(first_events.len(), 1);
6308 let OrderEventAny::Filled(explicit) = &first_events[0] else {
6309 panic!("expected the valid explicit fill");
6310 };
6311 assert_eq!(explicit.last_qty, Quantity::from("4.0"));
6312 assert_eq!(
6313 explicit.commission,
6314 Some(Money::new(0.25, Currency::USDT()))
6315 );
6316 assert_eq!(working.status(), OrderStatus::PartiallyFilled);
6317
6318 let expected = Money::new(1.5, Currency::USDT());
6319 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
6320 let retry_events = manager.reconcile_order_with_fills(
6321 &working,
6322 &report,
6323 &[],
6324 Some(&instrument),
6325 &mut fill_queue,
6326 Some(&retry_client),
6327 );
6328
6329 assert_eq!(retry_events.len(), 1);
6330 let OrderEventAny::Filled(residual) = &retry_events[0] else {
6331 panic!("expected the inferred residual fill");
6332 };
6333 assert_eq!(residual.last_qty, Quantity::from("6.0"));
6334 assert_eq!(residual.last_px, Price::from("33.33"));
6335 assert_eq!(residual.commission, Some(expected));
6336 assert_eq!(
6337 *retry_client.seen.borrow(),
6338 Some((
6339 Quantity::from("6.0"),
6340 residual.last_px,
6341 residual.liquidity_side,
6342 )),
6343 "commission must use the exact price and liquidity carried by the residual fill"
6344 );
6345 }
6346
6347 #[rstest]
6348 fn test_cached_reconciliation_preserves_explicit_fill_side() {
6349 let (mut manager, _cache, order, mut report, instrument) = cached_commission_fixtures();
6350 report.order_status = OrderStatus::PartiallyFilled;
6351 report.filled_qty = Quantity::from("4.0");
6352 let trade_id = TradeId::from("T-CONFLICTING-SIDE");
6353 let explicit_fill = FillReport::new(
6354 report.account_id,
6355 report.instrument_id,
6356 report.venue_order_id,
6357 trade_id,
6358 OrderSide::Sell,
6359 Quantity::from("4.0"),
6360 Price::from("100.0"),
6361 Money::new(0.25, Currency::USDT()),
6362 LiquiditySide::Taker,
6363 report.client_order_id,
6364 None,
6365 UnixNanos::from(1),
6366 UnixNanos::from(1),
6367 None,
6368 );
6369 let mut fill_queue = ReconciliationFillQueue::default();
6370
6371 let events = manager.reconcile_order_with_fills(
6372 &order,
6373 &report,
6374 &[&explicit_fill],
6375 Some(&instrument),
6376 &mut fill_queue,
6377 None,
6378 );
6379
6380 assert_eq!(events.len(), 1);
6381 let OrderEventAny::Filled(fill) = &events[0] else {
6382 panic!("expected the explicit fill");
6383 };
6384 assert_eq!(fill.trade_id, trade_id);
6385 assert_eq!(fill.order_side, OrderSide::Sell);
6386 assert_eq!(fill.last_qty, Quantity::from("4.0"));
6387 assert_eq!(fill.last_px, Price::from("100.0"));
6388 }
6389
6390 #[rstest]
6391 fn test_cached_snapshot_without_instrument_preserves_terminal_transition() {
6392 let (mut manager, _cache, order, mut report, _instrument) = cached_commission_fixtures();
6393 report.order_status = OrderStatus::Canceled;
6394 let explicit_fill = FillReport::new(
6395 report.account_id,
6396 report.instrument_id,
6397 report.venue_order_id,
6398 TradeId::from("T-MISSING-INSTRUMENT"),
6399 OrderSide::Buy,
6400 Quantity::from("4.0"),
6401 Price::from("100.0"),
6402 Money::new(0.25, Currency::USDT()),
6403 LiquiditySide::Taker,
6404 report.client_order_id,
6405 None,
6406 UnixNanos::from(1),
6407 UnixNanos::from(1),
6408 None,
6409 );
6410 let mut fill_queue = ReconciliationFillQueue::default();
6411
6412 let events = manager.reconcile_order_with_fills(
6413 &order,
6414 &report,
6415 &[&explicit_fill],
6416 None,
6417 &mut fill_queue,
6418 None,
6419 );
6420
6421 assert_eq!(events.len(), 1);
6422 assert!(matches!(events[0], OrderEventAny::Canceled(_)));
6423 }
6424
6425 #[rstest]
6426 fn test_continuous_reconciliation_uses_source_client_and_retries_commission() {
6427 let (mut manager, _cache, order, report, _instrument) = cached_commission_fixtures();
6428 let client_id = ClientId::from("STUB");
6429 let check = OpenOrderReportCheck {
6430 command: GenerateOrderStatusReports::new(
6431 UUID4::new(),
6432 UnixNanos::from(1),
6433 true,
6434 None,
6435 None,
6436 None,
6437 None,
6438 None,
6439 ),
6440 filtered_orders: vec![order],
6441 client_coverage: IndexMap::from([(
6442 report.client_order_id.unwrap(),
6443 ReportClientCoverage::Resolved(IndexSet::from([client_id])),
6444 )]),
6445 start: None,
6446 };
6447 let queried_clients = IndexSet::from([client_id]);
6448 let failed_clients = IndexSet::new();
6449 let failing_client = CommissionStubClient::new(CommissionOutcome::Failure);
6450
6451 let failed = manager.reconcile_open_order_reports(
6452 &check,
6453 vec![SourcedOrderStatusReport {
6454 client_id,
6455 report: report.clone(),
6456 }],
6457 &queried_clients,
6458 &failed_clients,
6459 &[&failing_client],
6460 );
6461
6462 assert!(failed.events.is_empty());
6463
6464 let expected = Money::new(1.5, Currency::USDT());
6465 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
6466 let retry = manager.reconcile_open_order_reports(
6467 &check,
6468 vec![SourcedOrderStatusReport { client_id, report }],
6469 &queried_clients,
6470 &failed_clients,
6471 &[&retry_client],
6472 );
6473
6474 assert_eq!(retry.events.len(), 1);
6475 let OrderEventAny::Filled(fill) = &retry.events[0] else {
6476 panic!("expected inferred fill on valid retry");
6477 };
6478 assert_eq!(fill.last_qty, Quantity::from("10.0"));
6479 assert_eq!(fill.commission, Some(expected));
6480 }
6481
6482 #[rstest]
6483 fn test_open_check_lookback_exclusion_warns_once_without_reconciliation_actions() {
6484 let client_order_id = ClientOrderId::from("O-LOOKBACK-OLD");
6485 let venue_order_id = VenueOrderId::from("V-LOOKBACK-OLD");
6486 let client_id = ClientId::from("BINANCE");
6487 let instrument_id = crypto_perpetual_ethusdt().id();
6488 let clock = Rc::new(RefCell::new(TestClock::new()));
6489 let cache = Rc::new(RefCell::new(Cache::default()));
6490 insert_accepted_limit_order(
6491 &cache,
6492 client_order_id,
6493 venue_order_id,
6494 instrument_id,
6495 client_id,
6496 );
6497 let order = cache.borrow().order_owned(&client_order_id).unwrap();
6498 let cutoff = UnixNanos::from(order.ts_last().as_u64().saturating_add(1));
6499 let check = OpenOrderReportCheck {
6500 command: GenerateOrderStatusReports::new(
6501 UUID4::new(),
6502 UnixNanos::from(1),
6503 false,
6504 None,
6505 Some(cutoff),
6506 None,
6507 None,
6508 None,
6509 ),
6510 filtered_orders: vec![order],
6511 client_coverage: IndexMap::from([(
6512 client_order_id,
6513 ReportClientCoverage::Resolved(IndexSet::from([client_id])),
6514 )]),
6515 start: Some(cutoff),
6516 };
6517 let queried_clients = IndexSet::from([client_id]);
6518 let mut manager = ExecutionManager::new(
6519 clock,
6520 cache.clone(),
6521 ExecutionManagerConfig {
6522 open_check_open_only: false,
6523 ..Default::default()
6524 },
6525 )
6526 .expect("valid config");
6527
6528 for _ in 0..2 {
6529 let result = manager.reconcile_open_order_reports(
6530 &check,
6531 Vec::new(),
6532 &queried_clients,
6533 &IndexSet::new(),
6534 &[],
6535 );
6536
6537 assert!(result.events.is_empty());
6538 assert!(result.targeted_queries.is_empty());
6539 assert_eq!(
6540 cache.borrow().order(&client_order_id).unwrap().status(),
6541 OrderStatus::Accepted
6542 );
6543 assert!(!manager.recon_check_retries.contains_key(&client_order_id));
6544 assert!(!manager.order_query_recency.contains_key(&client_order_id));
6545 assert!(!manager.targeted_order_queries.contains(&client_order_id));
6546 assert_eq!(
6547 manager.open_check_lookback_warnings,
6548 IndexSet::from([client_order_id])
6549 );
6550 assert_eq!(manager.open_check_lookback_warnings.len(), 1);
6551 }
6552 }
6553
6554 #[rstest]
6555 fn test_open_check_lookback_warning_clears_at_boundary_and_rearms() {
6556 let client_order_id = ClientOrderId::from("O-LOOKBACK-REARM");
6557 let client_id = ClientId::from("BINANCE");
6558 let cache = Rc::new(RefCell::new(Cache::default()));
6559 insert_accepted_limit_order(
6560 &cache,
6561 client_order_id,
6562 VenueOrderId::from("V-LOOKBACK-REARM"),
6563 crypto_perpetual_ethusdt().id(),
6564 client_id,
6565 );
6566 let order = cache.borrow().order_owned(&client_order_id).unwrap();
6567 let old_cutoff = UnixNanos::from(order.ts_last().as_u64().saturating_add(1));
6568 let make_check = |start| OpenOrderReportCheck {
6569 command: GenerateOrderStatusReports::new(
6570 UUID4::new(),
6571 UnixNanos::from(1),
6572 false,
6573 None,
6574 Some(start),
6575 None,
6576 None,
6577 None,
6578 ),
6579 filtered_orders: vec![order.clone()],
6580 client_coverage: IndexMap::from([(
6581 client_order_id,
6582 ReportClientCoverage::Resolved(IndexSet::from([client_id])),
6583 )]),
6584 start: Some(start),
6585 };
6586 let queried_clients = IndexSet::new();
6587 let mut manager = ExecutionManager::new(
6588 Rc::new(RefCell::new(TestClock::new())),
6589 cache,
6590 ExecutionManagerConfig {
6591 open_check_open_only: false,
6592 ..Default::default()
6593 },
6594 )
6595 .expect("valid config");
6596
6597 manager.reconcile_open_order_reports(
6598 &make_check(old_cutoff),
6599 Vec::new(),
6600 &queried_clients,
6601 &IndexSet::new(),
6602 &[],
6603 );
6604 assert!(
6605 manager
6606 .open_check_lookback_warnings
6607 .contains(&client_order_id)
6608 );
6609
6610 let boundary = order.ts_last();
6611 let boundary_result = manager.reconcile_open_order_reports(
6612 &make_check(boundary),
6613 Vec::new(),
6614 &queried_clients,
6615 &IndexSet::new(),
6616 &[],
6617 );
6618 assert!(boundary_result.targeted_queries.is_empty());
6619 assert!(
6620 !manager
6621 .open_check_lookback_warnings
6622 .contains(&client_order_id)
6623 );
6624 assert!(
6625 manager
6626 .missing_order_coverage_warnings
6627 .contains(&client_order_id)
6628 );
6629
6630 manager.reconcile_open_order_reports(
6631 &make_check(old_cutoff),
6632 Vec::new(),
6633 &queried_clients,
6634 &IndexSet::new(),
6635 &[],
6636 );
6637 assert_eq!(
6638 manager.open_check_lookback_warnings,
6639 IndexSet::from([client_order_id])
6640 );
6641 }
6642
6643 #[rstest]
6644 fn test_venue_order_id_mapped_report_clears_old_order_lookback_warning() {
6645 let client_order_id = ClientOrderId::from("O-LOOKBACK-MAPPED");
6646 let venue_order_id = VenueOrderId::from("V-LOOKBACK-MAPPED");
6647 let client_id = ClientId::from("BINANCE");
6648 let instrument_id = crypto_perpetual_ethusdt().id();
6649 let cache = Rc::new(RefCell::new(Cache::default()));
6650 insert_accepted_limit_order(
6651 &cache,
6652 client_order_id,
6653 venue_order_id,
6654 instrument_id,
6655 client_id,
6656 );
6657 let order = cache.borrow().order_owned(&client_order_id).unwrap();
6658 let cutoff = UnixNanos::from(order.ts_last().as_u64().saturating_add(1));
6659 let check = OpenOrderReportCheck {
6660 command: GenerateOrderStatusReports::new(
6661 UUID4::new(),
6662 UnixNanos::from(1),
6663 false,
6664 None,
6665 Some(cutoff),
6666 None,
6667 None,
6668 None,
6669 ),
6670 filtered_orders: vec![order],
6671 client_coverage: IndexMap::from([(
6672 client_order_id,
6673 ReportClientCoverage::Resolved(IndexSet::from([client_id])),
6674 )]),
6675 start: Some(cutoff),
6676 };
6677 let mut manager = ExecutionManager::new(
6678 Rc::new(RefCell::new(TestClock::new())),
6679 cache,
6680 ExecutionManagerConfig {
6681 open_check_open_only: false,
6682 ..Default::default()
6683 },
6684 )
6685 .expect("valid config");
6686 manager.open_check_lookback_warnings.insert(client_order_id);
6687
6688 let report = OrderStatusReport::new(
6691 AccountId::from("TEST-001"),
6692 instrument_id,
6693 None,
6694 venue_order_id,
6695 OrderSide::Buy.into(),
6696 OrderType::Limit,
6697 TimeInForce::Gtc,
6698 OrderStatus::Accepted,
6699 Quantity::from("10.0"),
6700 Quantity::from("0.0"),
6701 UnixNanos::from(0),
6702 UnixNanos::from(0),
6703 UnixNanos::from(0),
6704 None,
6705 );
6706
6707 let result = manager.reconcile_open_order_reports(
6708 &check,
6709 vec![SourcedOrderStatusReport { client_id, report }],
6710 &IndexSet::from([client_id]),
6711 &IndexSet::new(),
6712 &[],
6713 );
6714
6715 assert!(result.targeted_queries.is_empty());
6716 assert!(
6717 !manager
6718 .open_check_lookback_warnings
6719 .contains(&client_order_id)
6720 );
6721 }
6722
6723 #[rstest]
6724 fn test_positive_report_clears_old_order_lookback_warning_without_reinserting_it() {
6725 let client_order_id = ClientOrderId::from("O-LOOKBACK-REPORTED");
6726 let venue_order_id = VenueOrderId::from("V-LOOKBACK-REPORTED");
6727 let client_id = ClientId::from("BINANCE");
6728 let instrument_id = crypto_perpetual_ethusdt().id();
6729 let cache = Rc::new(RefCell::new(Cache::default()));
6730 insert_accepted_limit_order(
6731 &cache,
6732 client_order_id,
6733 venue_order_id,
6734 instrument_id,
6735 client_id,
6736 );
6737 let order = cache.borrow().order_owned(&client_order_id).unwrap();
6738 let cutoff = UnixNanos::from(order.ts_last().as_u64().saturating_add(1));
6739 let check = OpenOrderReportCheck {
6740 command: GenerateOrderStatusReports::new(
6741 UUID4::new(),
6742 UnixNanos::from(1),
6743 false,
6744 None,
6745 Some(cutoff),
6746 None,
6747 None,
6748 None,
6749 ),
6750 filtered_orders: vec![order],
6751 client_coverage: IndexMap::from([(
6752 client_order_id,
6753 ReportClientCoverage::Resolved(IndexSet::from([client_id])),
6754 )]),
6755 start: Some(cutoff),
6756 };
6757 let mut manager = ExecutionManager::new(
6758 Rc::new(RefCell::new(TestClock::new())),
6759 cache,
6760 ExecutionManagerConfig {
6761 open_check_open_only: false,
6762 ..Default::default()
6763 },
6764 )
6765 .expect("valid config");
6766 manager.open_check_lookback_warnings.insert(client_order_id);
6767 let report = OrderStatusReport::new(
6768 AccountId::from("TEST-001"),
6769 instrument_id,
6770 Some(client_order_id),
6771 venue_order_id,
6772 OrderSide::Buy.into(),
6773 OrderType::Limit,
6774 TimeInForce::Gtc,
6775 OrderStatus::Accepted,
6776 Quantity::from("10.0"),
6777 Quantity::from("0.0"),
6778 UnixNanos::from(0),
6779 UnixNanos::from(0),
6780 UnixNanos::from(0),
6781 None,
6782 );
6783
6784 let result = manager.reconcile_open_order_reports(
6785 &check,
6786 vec![SourcedOrderStatusReport { client_id, report }],
6787 &IndexSet::from([client_id]),
6788 &IndexSet::new(),
6789 &[],
6790 );
6791
6792 assert!(result.targeted_queries.is_empty());
6793 assert!(
6794 !manager
6795 .open_check_lookback_warnings
6796 .contains(&client_order_id)
6797 );
6798 }
6799
6800 #[rstest]
6801 fn test_targeted_reconciliation_uses_source_client_and_retries_commission() {
6802 let (mut manager, _cache, _order, report, _instrument) = cached_commission_fixtures();
6803 let client_order_id = report.client_order_id.unwrap();
6804 let client_id = ClientId::from("STUB");
6805 let failing_client = CommissionStubClient::new(CommissionOutcome::Failure);
6806
6807 let failed = manager.reconcile_targeted_order_reports(
6808 vec![TargetedOrderReportResult {
6809 client_order_id,
6810 client_id: Some(client_id),
6811 report: Some(report.clone()),
6812 coverage_complete: true,
6813 }],
6814 &[&failing_client],
6815 );
6816
6817 assert!(failed.is_empty());
6818
6819 let expected = Money::new(1.5, Currency::USDT());
6820 let retry_client = CommissionStubClient::new(CommissionOutcome::Value(expected));
6821 let retry = manager.reconcile_targeted_order_reports(
6822 vec![TargetedOrderReportResult {
6823 client_order_id,
6824 client_id: Some(client_id),
6825 report: Some(report),
6826 coverage_complete: true,
6827 }],
6828 &[&retry_client],
6829 );
6830
6831 assert_eq!(retry.len(), 1);
6832 let OrderEventAny::Filled(fill) = &retry[0] else {
6833 panic!("expected inferred fill on valid targeted retry");
6834 };
6835 assert_eq!(fill.last_qty, Quantity::from("10.0"));
6836 assert_eq!(fill.commission, Some(expected));
6837 }
6838
6839 #[rstest]
6840 fn test_clear_recon_tracking_removes_targeted_query() {
6841 let clock = Rc::new(RefCell::new(TestClock::new()));
6842 let cache = Rc::new(RefCell::new(Cache::default()));
6843 let mut manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default())
6844 .expect("valid config");
6845 let client_order_id = ClientOrderId::from("O-TARGETED-CLEAR");
6846 manager.targeted_order_queries.insert(client_order_id);
6847 manager.open_check_lookback_warnings.insert(client_order_id);
6848
6849 manager.clear_recon_tracking(&client_order_id, true);
6850
6851 assert!(manager.targeted_order_queries.is_empty());
6852 assert!(manager.open_check_lookback_warnings.is_empty());
6853 }
6854
6855 #[rstest]
6856 fn test_register_inflight_skips_filtered_order() {
6857 let client_order_id = ClientOrderId::from("O-FILTERED-REGISTER");
6858 let clock = Rc::new(RefCell::new(TestClock::new()));
6859 let cache = Rc::new(RefCell::new(Cache::default()));
6860 let mut manager = ExecutionManager::new(
6861 clock,
6862 cache,
6863 ExecutionManagerConfig {
6864 filtered_client_order_ids: IndexSet::from([client_order_id]),
6865 ..Default::default()
6866 },
6867 )
6868 .expect("valid config");
6869
6870 manager.register_inflight(client_order_id);
6871
6872 assert!(!manager.inflight_checks.contains_key(&client_order_id));
6873 assert!(!manager.recon_check_retries.contains_key(&client_order_id));
6874 }
6875
6876 #[rstest]
6877 #[cfg_attr(
6878 not(all(feature = "simulation", madsim)),
6879 tokio::test(start_paused = true)
6880 )]
6881 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
6882 async fn test_inflight_check_retires_order_filtered_after_registration() {
6883 let client_order_id = ClientOrderId::from("O-FILTERED-LATE");
6884 let clock = Rc::new(RefCell::new(TestClock::new()));
6885 let cache = Rc::new(RefCell::new(Cache::default()));
6886 let mut manager = ExecutionManager::new(
6887 clock,
6888 cache,
6889 ExecutionManagerConfig {
6890 inflight_threshold_ms: 100,
6891 ..Default::default()
6892 },
6893 )
6894 .expect("valid config");
6895 manager.register_inflight(client_order_id);
6896 manager
6897 .config
6898 .filtered_client_order_ids
6899 .insert(client_order_id);
6900 dst::time::sleep(Duration::from_millis(101)).await;
6901
6902 let first = manager.check_inflight_orders();
6903
6904 assert!(first.events.is_empty());
6905 assert!(first.queries.is_empty());
6906 assert!(!manager.inflight_checks.contains_key(&client_order_id));
6907 assert!(!manager.recon_check_retries.contains_key(&client_order_id));
6908
6909 dst::time::sleep(Duration::from_millis(101)).await;
6910 let second = manager.check_inflight_orders();
6911 assert!(second.events.is_empty());
6912 assert!(second.queries.is_empty());
6913 assert!(!manager.inflight_checks.contains_key(&client_order_id));
6914 }
6915
6916 #[rstest]
6917 #[case(false, OrderStatus::PendingUpdate, true, true, true)]
6918 #[case(false, OrderStatus::Accepted, false, true, true)]
6919 #[case(false, OrderStatus::Canceled, false, true, false)]
6920 #[case(true, OrderStatus::PendingCancel, true, true, true)]
6921 #[case(true, OrderStatus::Accepted, false, true, true)]
6922 #[case(true, OrderStatus::Filled, false, true, false)]
6923 fn test_observe_order_status_report_tracking_matrix(
6924 #[case] with_fills: bool,
6925 #[case] status: OrderStatus,
6926 #[case] expect_inflight: bool,
6927 #[case] expect_activity: bool,
6928 #[case] expect_last_query: bool,
6929 ) {
6930 let client_order_id = ClientOrderId::from("O-STATUS-MATRIX");
6931 let clock = Rc::new(RefCell::new(TestClock::new()));
6932 let cache = Rc::new(RefCell::new(Cache::default()));
6933 let mut manager = ExecutionManager::new(clock, cache, ExecutionManagerConfig::default())
6934 .expect("valid config");
6935 manager.register_inflight(client_order_id);
6936 manager.order_query_recency.mark(client_order_id);
6937 manager
6938 .missing_order_coverage_warnings
6939 .insert(client_order_id);
6940 manager.unresolved_order_coverage.insert(client_order_id);
6941 manager.targeted_order_queries.insert(client_order_id);
6942 let order_report = OrderStatusReport::new(
6943 AccountId::from("TEST-001"),
6944 crypto_perpetual_ethusdt().id(),
6945 Some(client_order_id),
6946 VenueOrderId::from("V-STATUS-MATRIX"),
6947 OrderSide::Buy.into(),
6948 OrderType::Limit,
6949 TimeInForce::Gtc,
6950 status,
6951 Quantity::from("10.0"),
6952 Quantity::from("0.0"),
6953 UnixNanos::from(1_000),
6954 UnixNanos::from(1_000),
6955 UnixNanos::from(1_000),
6956 None,
6957 );
6958 let report = if with_fills {
6959 ExecutionReport::OrderWithFills(Box::new(order_report), Vec::new())
6960 } else {
6961 ExecutionReport::Order(Box::new(order_report))
6962 };
6963
6964 manager.observe_execution_report(&report);
6965
6966 assert_eq!(
6967 manager.inflight_checks.contains_key(&client_order_id),
6968 expect_inflight,
6969 );
6970 assert_eq!(
6971 manager.recon_check_retries.contains_key(&client_order_id),
6972 expect_inflight,
6973 );
6974 assert_eq!(
6975 manager.order_local_activity.contains_key(&client_order_id),
6976 expect_activity,
6977 );
6978 assert_eq!(
6979 manager.order_query_recency.contains_key(&client_order_id),
6980 expect_last_query,
6981 );
6982 assert_eq!(
6983 manager
6984 .missing_order_coverage_warnings
6985 .contains(&client_order_id),
6986 expect_inflight,
6987 );
6988 assert_eq!(
6989 manager.unresolved_order_coverage.contains(&client_order_id),
6990 expect_inflight,
6991 );
6992 assert_eq!(
6993 manager.targeted_order_queries.contains(&client_order_id),
6994 expect_inflight,
6995 );
6996 }
6997
6998 #[rstest]
6999 #[case(OrderStatus::PendingUpdate)]
7000 #[case(OrderStatus::PendingCancel)]
7001 fn test_accepted_report_during_pending_command_preserves_inflight_tracking(
7002 #[case] pending_status: OrderStatus,
7003 ) {
7004 let client_order_id = ClientOrderId::from("O-PENDING-COMMAND");
7005 let venue_order_id = VenueOrderId::from("V-PENDING-COMMAND");
7006 let account_id = AccountId::from("TEST-001");
7007 let client_id = ClientId::from("TEST");
7008 let instrument_id = crypto_perpetual_ethusdt().id();
7009 let clock = Rc::new(RefCell::new(TestClock::new()));
7010 let cache = Rc::new(RefCell::new(Cache::default()));
7011 insert_accepted_limit_order(
7012 &cache,
7013 client_order_id,
7014 venue_order_id,
7015 instrument_id,
7016 client_id,
7017 );
7018
7019 let order = cache.borrow().order_owned(&client_order_id).unwrap();
7020 let event = match pending_status {
7021 OrderStatus::PendingUpdate => OrderEventAny::PendingUpdate(
7022 OrderPendingUpdateSpec::builder()
7023 .trader_id(order.trader_id())
7024 .strategy_id(order.strategy_id())
7025 .instrument_id(instrument_id)
7026 .client_order_id(client_order_id)
7027 .account_id(account_id)
7028 .venue_order_id(venue_order_id)
7029 .build(),
7030 ),
7031 OrderStatus::PendingCancel => OrderEventAny::PendingCancel(
7032 OrderPendingCancelSpec::builder()
7033 .trader_id(order.trader_id())
7034 .strategy_id(order.strategy_id())
7035 .instrument_id(instrument_id)
7036 .client_order_id(client_order_id)
7037 .account_id(account_id)
7038 .venue_order_id(venue_order_id)
7039 .build(),
7040 ),
7041 _ => unreachable!(),
7042 };
7043 cache.borrow_mut().update_order(&event).unwrap();
7044
7045 let mut manager =
7046 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default())
7047 .expect("valid config");
7048 manager.register_inflight(client_order_id);
7049 manager.order_query_recency.mark(client_order_id);
7050 manager
7051 .missing_order_coverage_warnings
7052 .insert(client_order_id);
7053 manager.unresolved_order_coverage.insert(client_order_id);
7054 manager.targeted_order_queries.insert(client_order_id);
7055 let report = OrderStatusReport::new(
7056 account_id,
7057 instrument_id,
7058 Some(client_order_id),
7059 venue_order_id,
7060 OrderSide::Buy.into(),
7061 OrderType::Limit,
7062 TimeInForce::Gtc,
7063 OrderStatus::Accepted,
7064 Quantity::from("10.0"),
7065 Quantity::from("0.0"),
7066 UnixNanos::from(1_000),
7067 UnixNanos::from(1_000),
7068 UnixNanos::from(1_000),
7069 None,
7070 )
7071 .with_price(Price::from("100.0"));
7072
7073 manager.observe_execution_report(&ExecutionReport::Order(Box::new(report.clone())));
7074 let order = cache.borrow().order_owned(&client_order_id).unwrap();
7075 let events =
7076 generate_reconciliation_order_events(&order, &report, None, UnixNanos::from(1_000));
7077
7078 assert!(events.is_empty());
7079 assert_eq!(order.status(), pending_status);
7080 assert!(manager.inflight_checks.contains_key(&client_order_id));
7081 assert!(manager.recon_check_retries.contains_key(&client_order_id));
7082 assert!(manager.order_query_recency.contains_key(&client_order_id));
7083 assert!(manager.order_local_activity.contains_key(&client_order_id));
7084 assert!(
7085 manager
7086 .missing_order_coverage_warnings
7087 .contains(&client_order_id)
7088 );
7089 assert!(manager.unresolved_order_coverage.contains(&client_order_id));
7090 assert!(manager.targeted_order_queries.contains(&client_order_id));
7091 }
7092
7093 #[rstest]
7094 fn test_superseded_cancel_report_preserves_missing_order_grace() {
7095 let client_order_id = ClientOrderId::from("O-CANCEL-REPLACE");
7096 let old_venue_order_id = VenueOrderId::from("V-CANCEL-REPLACE-OLD");
7097 let new_venue_order_id = VenueOrderId::from("V-CANCEL-REPLACE-NEW");
7098 let account_id = AccountId::from("TEST-001");
7099 let client_id = ClientId::from("TEST");
7100 let instrument_id = crypto_perpetual_ethusdt().id();
7101 let clock = Rc::new(RefCell::new(TestClock::new()));
7102 let cache = Rc::new(RefCell::new(Cache::default()));
7103 insert_accepted_limit_order(
7104 &cache,
7105 client_order_id,
7106 old_venue_order_id,
7107 instrument_id,
7108 client_id,
7109 );
7110
7111 let order = cache.borrow().order_owned(&client_order_id).unwrap();
7112 let pending_update = OrderPendingUpdateSpec::builder()
7113 .trader_id(order.trader_id())
7114 .strategy_id(order.strategy_id())
7115 .instrument_id(order.instrument_id())
7116 .client_order_id(client_order_id)
7117 .account_id(account_id)
7118 .venue_order_id(old_venue_order_id)
7119 .build();
7120 cache
7121 .borrow_mut()
7122 .update_order(&OrderEventAny::PendingUpdate(pending_update))
7123 .unwrap();
7124 let order = cache.borrow().order_owned(&client_order_id).unwrap();
7125 let updated = OrderUpdatedSpec::builder()
7126 .trader_id(order.trader_id())
7127 .strategy_id(order.strategy_id())
7128 .instrument_id(order.instrument_id())
7129 .client_order_id(client_order_id)
7130 .quantity(order.quantity())
7131 .venue_order_id(new_venue_order_id)
7132 .account_id(account_id)
7133 .build();
7134 cache
7135 .borrow_mut()
7136 .update_order(&OrderEventAny::Updated(updated))
7137 .unwrap();
7138
7139 let mut manager = ExecutionManager::new(
7140 clock,
7141 cache.clone(),
7142 ExecutionManagerConfig {
7143 open_check_missing_retries: 1,
7144 ..Default::default()
7145 },
7146 )
7147 .expect("valid config");
7148 manager.record_local_activity(client_order_id);
7149 assert!(
7150 manager
7151 .prepare_missing_order_query(client_order_id)
7152 .is_none()
7153 );
7154
7155 let report = OrderStatusReport::new(
7156 account_id,
7157 instrument_id,
7158 Some(client_order_id),
7159 old_venue_order_id,
7160 OrderSide::Buy.into(),
7161 OrderType::Limit,
7162 TimeInForce::Gtc,
7163 OrderStatus::Canceled,
7164 Quantity::from("10.0"),
7165 Quantity::from("0.0"),
7166 UnixNanos::from(1_000),
7167 UnixNanos::from(2_000),
7168 UnixNanos::from(3_000),
7169 None,
7170 );
7171
7172 manager.observe_execution_report(&ExecutionReport::Order(Box::new(report.clone())));
7173 let order = cache.borrow().order_owned(&client_order_id).unwrap();
7174 let events =
7175 generate_reconciliation_order_events(&order, &report, None, UnixNanos::from(1_000));
7176
7177 assert!(events.is_empty());
7178 assert_eq!(order.status(), OrderStatus::Accepted);
7179 assert_eq!(order.venue_order_id(), Some(new_venue_order_id));
7180 assert!(manager.order_local_activity.contains_key(&client_order_id));
7181 assert!(
7182 manager
7183 .prepare_missing_order_query(client_order_id)
7184 .is_none()
7185 );
7186 assert_eq!(manager.recon_check_retry_count(&client_order_id), 0);
7187 }
7188
7189 #[rstest]
7190 #[cfg_attr(
7191 not(all(feature = "simulation", madsim)),
7192 tokio::test(start_paused = true)
7193 )]
7194 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
7195 async fn test_prune_order_local_activity_uses_open_check_threshold() {
7196 let old_id = ClientOrderId::from("O-ACTIVITY-OLD");
7197 let fresh_id = ClientOrderId::from("O-ACTIVITY-FRESH");
7198 let clock = Rc::new(RefCell::new(TestClock::new()));
7199 let cache = Rc::new(RefCell::new(Cache::default()));
7200 let mut manager = ExecutionManager::new(
7201 clock,
7202 cache,
7203 ExecutionManagerConfig {
7204 open_check_threshold_ns: 100_000_000,
7205 ..Default::default()
7206 },
7207 )
7208 .expect("valid config");
7209 manager.record_local_activity(old_id);
7210 dst::time::sleep(Duration::from_millis(101)).await;
7211 manager.record_local_activity(fresh_id);
7212
7213 manager.prune_order_local_activity();
7214
7215 assert!(!manager.order_local_activity.contains_key(&old_id));
7216 assert!(manager.order_local_activity.contains_key(&fresh_id));
7217 }
7218
7219 #[rstest]
7220 fn test_prepare_open_order_report_check_builds_bulk_command_with_config() {
7221 let lookback_mins = 5_u64;
7222 let lookback_ns = lookback_mins * 60 * NANOSECONDS_IN_SECOND;
7223 let clock = Rc::new(RefCell::new(TestClock::new()));
7224 let cache = Rc::new(RefCell::new(Cache::default()));
7225 let mut manager = ExecutionManager::new(
7226 clock.clone(),
7227 cache.clone(),
7228 ExecutionManagerConfig {
7229 open_check_lookback_mins: Some(lookback_mins),
7230 open_check_open_only: false,
7231 reconciliation_instrument_ids: IndexSet::from([crypto_perpetual_ethusdt().id()]),
7232 ..Default::default()
7233 },
7234 )
7235 .expect("valid config");
7236 let included_id = ClientOrderId::from("O-REPORT-001");
7237 let excluded_id = ClientOrderId::from("O-REPORT-002");
7238 let included_instrument_id = crypto_perpetual_ethusdt().id();
7239 let excluded_instrument_id = xbtusd_bitmex().id();
7240
7241 cache
7242 .borrow_mut()
7243 .add_instrument(InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt()))
7244 .unwrap();
7245 cache
7246 .borrow_mut()
7247 .add_instrument(InstrumentAny::CryptoPerpetual(xbtusd_bitmex()))
7248 .unwrap();
7249 insert_accepted_limit_order(
7250 &cache,
7251 included_id,
7252 VenueOrderId::from("V-REPORT-001"),
7253 included_instrument_id,
7254 ClientId::from("BINANCE"),
7255 );
7256 insert_accepted_limit_order(
7257 &cache,
7258 excluded_id,
7259 VenueOrderId::from("V-REPORT-002"),
7260 excluded_instrument_id,
7261 ClientId::from("BITMEX"),
7262 );
7263 clock
7264 .borrow_mut()
7265 .advance_time(UnixNanos::from(lookback_ns * 2), true);
7266
7267 let ts_now = clock.borrow().timestamp_ns();
7268 let command_id = UUID4::new();
7269 let check = manager.prepare_open_order_report_check(command_id, &[]);
7270
7271 assert_eq!(check.command.command_id, command_id);
7272 assert_eq!(check.command.ts_init, ts_now);
7273 assert!(!check.command.open_only);
7274 assert_eq!(check.command.instrument_id, None);
7275 assert_eq!(
7276 check.command.start,
7277 Some(ts_now.saturating_sub_ns(lookback_ns))
7278 );
7279 assert_eq!(check.command.end, None);
7280 assert_eq!(check.command.log_receipt_level, LogLevel::Debug);
7281 assert_eq!(check.start, check.command.start);
7282 assert_eq!(check.filtered_orders.len(), 1);
7283 assert_eq!(check.filtered_orders[0].client_order_id(), included_id);
7284 }
7285
7286 #[rstest]
7287 fn test_prepare_position_report_check_builds_bulk_command_with_coverage() {
7288 let clock = Rc::new(RefCell::new(TestClock::new()));
7289 let cache = Rc::new(RefCell::new(Cache::default()));
7290 let manager = ExecutionManager::new(
7291 clock.clone(),
7292 cache.clone(),
7293 ExecutionManagerConfig {
7294 reconciliation_instrument_ids: IndexSet::from([crypto_perpetual_ethusdt().id()]),
7295 ..Default::default()
7296 },
7297 )
7298 .expect("valid config");
7299 let included_instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
7300 let excluded_instrument = InstrumentAny::CryptoPerpetual(xbtusd_bitmex());
7301
7302 cache
7303 .borrow_mut()
7304 .add_instrument(included_instrument.clone())
7305 .unwrap();
7306 cache
7307 .borrow_mut()
7308 .add_instrument(excluded_instrument.clone())
7309 .unwrap();
7310 let included_position = insert_open_position(
7311 &cache,
7312 &included_instrument,
7313 PositionId::from("P-REPORT-001"),
7314 OrderSide::Buy,
7315 "5.0",
7316 "3000.00",
7317 );
7318 insert_open_position(
7319 &cache,
7320 &excluded_instrument,
7321 PositionId::from("P-REPORT-002"),
7322 OrderSide::Buy,
7323 "2.0",
7324 "40000.00",
7325 );
7326
7327 let ts_now = clock.borrow().timestamp_ns();
7328 let command_id = UUID4::new();
7329 let check = manager.prepare_position_report_check(command_id, &[]);
7330 let key = (
7331 included_position.instrument_id,
7332 included_position.account_id,
7333 );
7334
7335 assert_eq!(check.command.command_id, command_id);
7336 assert_eq!(check.command.ts_init, ts_now);
7337 assert_eq!(check.command.instrument_id, None);
7338 assert_eq!(check.command.start, None);
7339 assert_eq!(check.command.end, None);
7340 assert_eq!(check.command.log_receipt_level, LogLevel::Debug);
7341 assert_eq!(check.client_coverage.len(), 1);
7342 assert_eq!(
7343 check.client_coverage.get(&key),
7344 Some(&ReportClientCoverage::Unresolved)
7345 );
7346 assert_eq!(check.activity_revisions.get(&key), Some(&0));
7347 }
7348
7349 #[rstest]
7350 #[cfg(feature = "node")]
7351 fn test_prepare_position_fill_report_plan_uses_configured_lookback() {
7352 let lookback_mins = 7_u64;
7353 let lookback_ns = lookback_mins * 60 * NANOSECONDS_IN_SECOND;
7354 let clock = Rc::new(RefCell::new(TestClock::new()));
7355 let cache = Rc::new(RefCell::new(Cache::default()));
7356 let mut manager = ExecutionManager::new(
7357 clock.clone(),
7358 cache.clone(),
7359 ExecutionManagerConfig {
7360 position_check_lookback_mins: lookback_mins,
7361 position_check_threshold_ns: 0,
7362 ..Default::default()
7363 },
7364 )
7365 .expect("valid config");
7366 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
7367 cache
7368 .borrow_mut()
7369 .add_instrument(instrument.clone())
7370 .unwrap();
7371 let position = insert_open_position(
7372 &cache,
7373 &instrument,
7374 PositionId::from("P-FILL-LOOKBACK"),
7375 OrderSide::Buy,
7376 "1.0",
7377 "3000.00",
7378 );
7379 clock
7380 .borrow_mut()
7381 .advance_time(UnixNanos::from(lookback_ns * 2), true);
7382 let client = PositionCoverageStubClient;
7383 let clients: [&dyn ExecutionClient; 1] = [&client];
7384 let mut check = manager.prepare_position_report_check(UUID4::new(), &clients);
7385 let query_end = clock.borrow().timestamp_ns();
7386 let report = PositionStatusReport::new(
7387 position.account_id,
7388 position.instrument_id,
7389 PositionSide::Long,
7390 Quantity::from("2.0"),
7391 query_end,
7392 query_end,
7393 None,
7394 None,
7395 Some(dec!(3000.00)),
7396 );
7397 let queried_clients = IndexSet::from([client.client_id()]);
7398
7399 let plan = manager.prepare_position_fill_report_plan(
7400 &mut check,
7401 &[report],
7402 &queried_clients,
7403 &IndexSet::new(),
7404 &clients,
7405 );
7406
7407 assert_eq!(
7408 plan.discrepancy_keys,
7409 IndexSet::from([(position.instrument_id, position.account_id)])
7410 );
7411 assert_eq!(plan.queries.len(), 1);
7412 let query = &plan.queries[0];
7413 assert_eq!(
7414 (query.key, query.client_id),
7415 (
7416 (position.instrument_id, position.account_id),
7417 client.client_id()
7418 )
7419 );
7420 assert_eq!(query.command.instrument_id, Some(position.instrument_id));
7421 assert_eq!(query.command.venue_order_id, None);
7422 assert_eq!(
7423 query.command.start,
7424 Some(query_end.saturating_sub_ns(lookback_ns))
7425 );
7426 assert_eq!(query.command.end, Some(query_end));
7427 assert_eq!(query.command.correlation_id, Some(check.command.command_id));
7428 assert_eq!(query.command.log_receipt_level, LogLevel::Debug);
7429 }
7430
7431 #[rstest]
7432 #[cfg(feature = "node")]
7433 fn test_prepare_position_fill_report_plan_defers_position_opened_during_request() {
7434 let clock = Rc::new(RefCell::new(TestClock::new()));
7435 let cache = Rc::new(RefCell::new(Cache::default()));
7436 let mut manager = ExecutionManager::new(
7437 clock,
7438 cache.clone(),
7439 ExecutionManagerConfig {
7440 position_check_threshold_ns: 0,
7441 ..Default::default()
7442 },
7443 )
7444 .expect("valid config");
7445 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
7446 cache
7447 .borrow_mut()
7448 .add_instrument(instrument.clone())
7449 .unwrap();
7450 let client = PositionCoverageStubClient;
7451 let clients: [&dyn ExecutionClient; 1] = [&client];
7452 let mut check = manager.prepare_position_report_check(UUID4::new(), &clients);
7453 let position = insert_open_position(
7454 &cache,
7455 &instrument,
7456 PositionId::from("P-FILL-DURING-REQUEST"),
7457 OrderSide::Buy,
7458 "1.0",
7459 "3000.00",
7460 );
7461 manager.record_position_activity(position.instrument_id, position.account_id);
7462 let report = PositionStatusReport::new(
7463 position.account_id,
7464 position.instrument_id,
7465 PositionSide::Long,
7466 Quantity::from("2.0"),
7467 UnixNanos::from(1_000_000),
7468 UnixNanos::from(1_000_000),
7469 None,
7470 None,
7471 Some(dec!(3000.00)),
7472 );
7473
7474 let plan = manager.prepare_position_fill_report_plan(
7475 &mut check,
7476 &[report],
7477 &IndexSet::from([client.client_id()]),
7478 &IndexSet::new(),
7479 &clients,
7480 );
7481
7482 assert_eq!(
7483 plan.discrepancy_keys,
7484 IndexSet::from([(position.instrument_id, position.account_id)])
7485 );
7486 assert!(plan.queries.is_empty());
7487 assert_eq!(
7488 check
7489 .activity_revisions
7490 .get(&(position.instrument_id, position.account_id)),
7491 Some(&0)
7492 );
7493 }
7494
7495 #[rstest]
7496 #[cfg(feature = "node")]
7497 fn test_prepare_position_report_check_uses_live_client_bulk_coverage() {
7498 let clock = Rc::new(RefCell::new(TestClock::new()));
7499 let cache = Rc::new(RefCell::new(Cache::default()));
7500 let manager =
7501 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default())
7502 .expect("valid config");
7503 let derivative = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
7504 let spot = test_bybit_spot_instrument();
7505 cache
7506 .borrow_mut()
7507 .add_instrument(derivative.clone())
7508 .unwrap();
7509 cache.borrow_mut().add_instrument(spot.clone()).unwrap();
7510 let derivative_position = insert_open_position(
7511 &cache,
7512 &derivative,
7513 PositionId::from("P-DERIVATIVE-COVERAGE"),
7514 OrderSide::Buy,
7515 "5.0",
7516 "3000.00",
7517 );
7518 let spot_position = insert_open_position(
7519 &cache,
7520 &spot,
7521 PositionId::from("P-SPOT-COVERAGE"),
7522 OrderSide::Buy,
7523 "2.0",
7524 "2000.00",
7525 );
7526 let client = LiveExecutionClient::new(Box::new(PositionCoverageStubClient));
7527 let client: &dyn ExecutionClient = &client;
7528
7529 let check = manager.prepare_position_report_check(UUID4::new(), &[client]);
7530 let client_id = ClientId::from("BYBIT");
7531
7532 assert_eq!(
7533 check.client_coverage.get(&(
7534 derivative_position.instrument_id,
7535 derivative_position.account_id
7536 )),
7537 Some(&ReportClientCoverage::Resolved(IndexSet::from([client_id])))
7538 );
7539 assert_eq!(
7540 check
7541 .client_coverage
7542 .get(&(spot_position.instrument_id, spot_position.account_id)),
7543 Some(&ReportClientCoverage::Unavailable(IndexSet::from([
7544 client_id
7545 ])))
7546 );
7547 }
7548
7549 #[rstest]
7550 fn test_position_reconciliation_preserves_unavailable_spot_coverage() {
7551 let clock = Rc::new(RefCell::new(TestClock::new()));
7552 let cache = Rc::new(RefCell::new(Cache::default()));
7553 let mut manager =
7554 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default())
7555 .expect("valid config");
7556 let derivative = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
7557 let spot = test_bybit_spot_instrument();
7558 cache
7559 .borrow_mut()
7560 .add_instrument(derivative.clone())
7561 .unwrap();
7562 cache.borrow_mut().add_instrument(spot.clone()).unwrap();
7563 let derivative_position = insert_open_position(
7564 &cache,
7565 &derivative,
7566 PositionId::from("P-DERIVATIVE-RECONCILE"),
7567 OrderSide::Buy,
7568 "5.0",
7569 "3000.00",
7570 );
7571 let spot_position = insert_open_position(
7572 &cache,
7573 &spot,
7574 PositionId::from("P-SPOT-PRESERVED"),
7575 OrderSide::Buy,
7576 "2.0",
7577 "2000.00",
7578 );
7579 let client = PositionCoverageStubClient;
7580 let check = manager.prepare_position_report_check(UUID4::new(), &[&client]);
7581 let queried_clients = IndexSet::from([client.client_id()]);
7582
7583 let events = manager.reconcile_position_reports(
7584 &check,
7585 Vec::new(),
7586 &queried_clients,
7587 &IndexSet::new(),
7588 );
7589
7590 assert!(events.iter().any(|event| {
7591 matches!(event, OrderEventAny::Filled(fill) if fill.instrument_id == derivative_position.instrument_id)
7592 }));
7593 assert!(!events.iter().any(|event| {
7594 matches!(event, OrderEventAny::Filled(fill) if fill.instrument_id == spot_position.instrument_id)
7595 }));
7596 }
7597
7598 #[rstest]
7599 fn test_position_reconciliation_preserves_spot_position_when_client_query_fails() {
7600 let clock = Rc::new(RefCell::new(TestClock::new()));
7601 let cache = Rc::new(RefCell::new(Cache::default()));
7602 let mut manager =
7603 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default())
7604 .expect("valid config");
7605 let instrument = test_bybit_spot_instrument();
7606 cache
7607 .borrow_mut()
7608 .add_instrument(instrument.clone())
7609 .unwrap();
7610 let position = insert_open_position(
7611 &cache,
7612 &instrument,
7613 PositionId::from("P-SPOT-QUERY-FAILED"),
7614 OrderSide::Buy,
7615 "5.0",
7616 "3000.00",
7617 );
7618 let key = (position.instrument_id, position.account_id);
7619 let client_id = ClientId::from("BYBIT");
7620 let mut check = manager.prepare_position_report_check(UUID4::new(), &[]);
7621 check.client_coverage.insert(
7622 key,
7623 ReportClientCoverage::Resolved(IndexSet::from([client_id])),
7624 );
7625 let queried_clients = IndexSet::from([client_id]);
7626 let failed_clients = IndexSet::from([client_id]);
7627
7628 let events = manager.reconcile_position_reports(
7629 &check,
7630 Vec::new(),
7631 &queried_clients,
7632 &failed_clients,
7633 );
7634
7635 assert!(
7636 !events
7637 .iter()
7638 .any(|event| matches!(event, OrderEventAny::Filled(_))),
7639 "a failed bulk query must not generate a synthetic closing fill",
7640 );
7641 let cached_position = cache.borrow().position(&position.id).unwrap().clone();
7642 assert!(cached_position.is_open());
7643 assert_eq!(cached_position.quantity, Quantity::from("5.0"));
7644 }
7645
7646 #[rstest]
7647 #[cfg_attr(
7648 not(all(feature = "simulation", madsim)),
7649 tokio::test(start_paused = true)
7650 )]
7651 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
7652 async fn test_position_report_check_defers_activity_recorded_during_delayed_request() {
7653 let clock = Rc::new(RefCell::new(TestClock::new()));
7654 let cache = Rc::new(RefCell::new(Cache::default()));
7655 let mut manager = ExecutionManager::new(
7656 clock,
7657 cache.clone(),
7658 ExecutionManagerConfig {
7659 position_check_threshold_ns: 5_000_000_000,
7660 ..Default::default()
7661 },
7662 )
7663 .expect("valid config");
7664 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
7665 let instrument_id = instrument.id();
7666 let position = insert_open_position(
7667 &cache,
7668 &instrument,
7669 PositionId::from("P-ACTIVITY-DURING-REQUEST"),
7670 OrderSide::Buy,
7671 "5.0",
7672 "3000.00",
7673 );
7674 cache
7675 .borrow_mut()
7676 .add_instrument(instrument.clone())
7677 .unwrap();
7678 let account_id = position.account_id;
7679 let check = manager.prepare_position_report_check(UUID4::new(), &[]);
7680 let report = PositionStatusReport::new(
7681 account_id,
7682 instrument_id,
7683 PositionSide::Long,
7684 Quantity::from("5.0"),
7685 UnixNanos::from(1_000_000),
7686 UnixNanos::from(1_000_000),
7687 None,
7688 None,
7689 Some(Decimal::from(3000)),
7690 );
7691
7692 let closed_position = close_long_position(
7693 position,
7694 &instrument,
7695 TradeId::from("T-ACTIVITY-DURING-REQUEST"),
7696 );
7697 cache
7698 .borrow_mut()
7699 .update_position(&closed_position)
7700 .unwrap();
7701 manager.record_position_activity(instrument_id, account_id);
7702
7703 dst::time::sleep(Duration::from_secs(6)).await;
7705
7706 let events = manager.reconcile_position_reports(
7707 &check,
7708 vec![report],
7709 &IndexSet::new(),
7710 &IndexSet::new(),
7711 );
7712
7713 assert!(
7714 !events.iter().any(|event| {
7715 matches!(
7716 event,
7717 OrderEventAny::Filled(fill)
7718 if fill.order_side == OrderSide::Buy
7719 && fill.last_qty == Quantity::from("5.0")
7720 )
7721 }),
7722 "activity recorded after the request started must defer A's stale report",
7723 );
7724 }
7725
7726 #[rstest]
7727 fn test_position_report_check_does_not_defer_activity_recorded_before_request() {
7728 let clock = Rc::new(RefCell::new(TestClock::new()));
7729 let cache = Rc::new(RefCell::new(Cache::default()));
7730 let mut manager = ExecutionManager::new(
7731 clock,
7732 cache.clone(),
7733 ExecutionManagerConfig {
7734 position_check_threshold_ns: 0,
7735 ..Default::default()
7736 },
7737 )
7738 .expect("valid config");
7739 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
7740 let instrument_id = instrument.id();
7741 let position = insert_open_position(
7742 &cache,
7743 &instrument,
7744 PositionId::from("P-ACTIVITY-BEFORE-REQUEST"),
7745 OrderSide::Buy,
7746 "5.0",
7747 "3000.00",
7748 );
7749 cache.borrow_mut().add_instrument(instrument).unwrap();
7750 let account_id = position.account_id;
7751 manager.record_position_activity(instrument_id, account_id);
7752 let check = manager.prepare_position_report_check(UUID4::new(), &[]);
7753 let report = PositionStatusReport::new(
7754 account_id,
7755 instrument_id,
7756 PositionSide::Long,
7757 Quantity::from("10.0"),
7758 UnixNanos::from(1_000_000),
7759 UnixNanos::from(1_000_000),
7760 None,
7761 None,
7762 Some(Decimal::from(3000)),
7763 );
7764
7765 let events = manager.reconcile_position_reports(
7766 &check,
7767 vec![report],
7768 &IndexSet::new(),
7769 &IndexSet::new(),
7770 );
7771
7772 let fills: Vec<_> = events
7773 .iter()
7774 .filter_map(|event| match event {
7775 OrderEventAny::Filled(fill) => Some(fill),
7776 _ => None,
7777 })
7778 .collect();
7779
7780 assert_eq!(fills.len(), 1);
7781 assert_eq!(fills[0].order_side, OrderSide::Buy);
7782 assert_eq!(fills[0].last_qty, Quantity::from("5.0"));
7783 assert_eq!(fills[0].commission, None);
7784 }
7785
7786 #[rstest]
7787 fn test_mass_status_projects_companion_fill_before_void_correction() {
7788 let clock = Rc::new(RefCell::new(TestClock::new()));
7789 let cache = Rc::new(RefCell::new(Cache::default()));
7790 let mut manager =
7791 ExecutionManager::new(clock, cache.clone(), ExecutionManagerConfig::default())
7792 .expect("valid config");
7793 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
7794 let client_order_id = ClientOrderId::from("O-MASS-VOID-001");
7795 let venue_order_id = VenueOrderId::from("V-MASS-VOID-001");
7796 let account_id = AccountId::from("TEST-001");
7797 cache
7798 .borrow_mut()
7799 .add_instrument(instrument.clone())
7800 .unwrap();
7801 insert_accepted_limit_order(
7802 &cache,
7803 client_order_id,
7804 venue_order_id,
7805 instrument.id(),
7806 ClientId::from("BINANCE"),
7807 );
7808 let order = cache.borrow().order_owned(&client_order_id).unwrap();
7809 let initial_fill = TestOrderEventStubs::filled(
7810 &order,
7811 &instrument,
7812 Some(TradeId::from("T-MASS-VOID-INITIAL")),
7813 None,
7814 Some(Price::from("100.0")),
7815 Some(Quantity::from("6.0")),
7816 Some(LiquiditySide::Taker),
7817 None,
7818 None,
7819 Some(account_id),
7820 );
7821 cache.borrow_mut().update_order(&initial_fill).unwrap();
7822 let order = cache.borrow().order_owned(&client_order_id).unwrap();
7823 let report = OrderStatusReport::new(
7824 account_id,
7825 instrument.id(),
7826 Some(client_order_id),
7827 venue_order_id,
7828 OrderSide::Buy.into(),
7829 OrderType::Limit,
7830 TimeInForce::Gtc,
7831 OrderStatus::Canceled,
7832 Quantity::from("10.0"),
7833 Quantity::from("5.0"),
7834 UnixNanos::from(1_000),
7835 UnixNanos::from(1_000),
7836 UnixNanos::from(1_000),
7837 None,
7838 );
7839 let companion_fill = FillReport::new(
7840 account_id,
7841 instrument.id(),
7842 venue_order_id,
7843 TradeId::from("T-MASS-VOID-COMPANION"),
7844 OrderSide::Buy,
7845 Quantity::from("1.0"),
7846 Price::from("100.0"),
7847 Money::zero(instrument.quote_currency()),
7848 LiquiditySide::Taker,
7849 Some(client_order_id),
7850 None,
7851 UnixNanos::from(900),
7852 UnixNanos::from(1_000),
7853 None,
7854 );
7855
7856 let mut fill_queue = ReconciliationFillQueue::default();
7857 let events = manager.reconcile_order_with_fills(
7858 &order,
7859 &report,
7860 &[&companion_fill],
7861 Some(&instrument),
7862 &mut fill_queue,
7863 None,
7864 );
7865 let mut projected = order;
7866 for event in &events {
7867 projected.apply(event.clone()).unwrap();
7868 }
7869
7870 assert!(matches!(events[0], OrderEventAny::Filled(_)));
7871 assert_eq!(
7872 events
7873 .iter()
7874 .filter(|event| matches!(event, OrderEventAny::FillVoided(_)))
7875 .count(),
7876 2
7877 );
7878 assert_eq!(projected.status(), OrderStatus::Canceled);
7879 assert_eq!(projected.filled_qty(), Quantity::from("5.0"));
7880 assert_eq!(projected.voided_qty(), Quantity::from("2.0"));
7881 }
7882
7883 fn insert_accepted_limit_order(
7884 cache: &Rc<RefCell<Cache>>,
7885 client_order_id: ClientOrderId,
7886 venue_order_id: VenueOrderId,
7887 instrument_id: InstrumentId,
7888 client_id: ClientId,
7889 ) {
7890 let account_id = AccountId::from("TEST-001");
7891 let order = OrderTestBuilder::new(OrderType::Limit)
7892 .client_order_id(client_order_id)
7893 .instrument_id(instrument_id)
7894 .quantity(Quantity::from("10.0"))
7895 .price(Price::from("100.0"))
7896 .build();
7897 let submitted = TestOrderEventStubs::submitted(&order, account_id);
7898 cache
7899 .borrow_mut()
7900 .add_order(order, None, Some(client_id), false)
7901 .unwrap();
7902 let order = cache.borrow_mut().update_order(&submitted).unwrap();
7903 let accepted = TestOrderEventStubs::accepted(&order, account_id, venue_order_id);
7904 cache.borrow_mut().update_order(&accepted).unwrap();
7905 }
7906
7907 fn test_bybit_spot_instrument() -> InstrumentAny {
7908 InstrumentAny::CurrencyPair(
7909 CurrencyPair::builder()
7910 .instrument_id(InstrumentId::from("ETHUSDT-SPOT.BYBIT"))
7911 .raw_symbol(Symbol::from("ETHUSDT"))
7912 .base_currency(Currency::from("ETH"))
7913 .quote_currency(Currency::from("USDT"))
7914 .price_precision(2)
7915 .size_precision(5)
7916 .price_increment(Price::from("0.01"))
7917 .size_increment(Quantity::from("0.00001"))
7918 .ts_event(UnixNanos::default())
7919 .ts_init(UnixNanos::default())
7920 .build()
7921 .unwrap(),
7922 )
7923 }
7924
7925 fn insert_open_position(
7926 cache: &Rc<RefCell<Cache>>,
7927 instrument: &InstrumentAny,
7928 position_id: PositionId,
7929 side: OrderSide,
7930 quantity: &str,
7931 price: &str,
7932 ) -> Position {
7933 let order = OrderTestBuilder::new(OrderType::Market)
7934 .instrument_id(instrument.id())
7935 .side(side)
7936 .quantity(Quantity::from(quantity))
7937 .build();
7938 let fill = TestOrderEventStubs::filled(
7939 &order,
7940 instrument,
7941 Some(TradeId::new("T-REPORT-001")),
7942 Some(position_id),
7943 Some(Price::from(price)),
7944 Some(Quantity::from(quantity)),
7945 None,
7946 None,
7947 None,
7948 Some(AccountId::from("TEST-001")),
7949 );
7950 let order_filled: OrderFilled = fill.into();
7951 let position = Position::new(instrument, order_filled);
7952 cache
7953 .borrow_mut()
7954 .add_position(&position, OmsType::Hedging)
7955 .unwrap();
7956 position
7957 }
7958
7959 fn close_long_position(
7960 mut position: Position,
7961 instrument: &InstrumentAny,
7962 trade_id: TradeId,
7963 ) -> Position {
7964 let order = OrderTestBuilder::new(OrderType::Market)
7965 .instrument_id(instrument.id())
7966 .side(OrderSide::Sell)
7967 .quantity(position.quantity)
7968 .build();
7969 let fill = TestOrderEventStubs::filled(
7970 &order,
7971 instrument,
7972 Some(trade_id),
7973 Some(position.id),
7974 Some(Price::from("3000.00")),
7975 Some(position.quantity),
7976 None,
7977 None,
7978 None,
7979 Some(position.account_id),
7980 );
7981 let order_filled: OrderFilled = fill.into();
7982 position.apply(&order_filled);
7983 position
7984 }
7985}