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