Skip to main content

nautilus_live/execution/
manager.rs

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