Skip to main content

futu_cache/trd_cache/
backend_merge.rs

1use super::{CachedOrder, OrderRelationSourceToken, TrdCache};
2use futu_domain_trade_account::{
3    BackendOrderMergeOrderFacts, plan_backend_authoritative_order_merge_like_cpp,
4};
5use futu_domain_trade_order::{
6    OrderSnapshotReplacementPlan, OrderSnapshotRowFacts, TradeOrderFacts,
7    plan_order_snapshot_replacement_like_cpp,
8};
9
10impl TrdCache {
11    /// v1.4.90 cache saga fix: merge backend authoritative orders while retaining fresh local stubs.
12    pub fn merge_preserving_stubs(&self, acc_id: u64, backend_orders: Vec<CachedOrder>) {
13        let _ = self.replace_order_snapshot_like_cpp(acc_id, &backend_orders);
14    }
15
16    /// Injectable-clock variant used by regression tests for fresh/stale stub boundaries.
17    pub fn merge_preserving_stubs_with_now(
18        &self,
19        acc_id: u64,
20        backend_orders: Vec<CachedOrder>,
21        now_ms: u64,
22    ) {
23        let _ = self.replace_order_snapshot_like_cpp_with_now(acc_id, &backend_orders, now_ms);
24    }
25
26    /// Atomically plan C++ `UpdateOrderList` notifications and replace the
27    /// backend-authoritative snapshot while retaining Rust local compatibility
28    /// rows (fresh optimistic stubs and C++ local Deleted rows).
29    pub fn replace_order_snapshot_like_cpp(
30        &self,
31        acc_id: u64,
32        backend_orders: &[CachedOrder],
33    ) -> OrderSnapshotReplacementPlan {
34        self.replace_order_snapshot_like_cpp_internal(
35            acc_id,
36            backend_orders,
37            Self::now_ms(),
38            None,
39            true,
40        )
41    }
42
43    /// Production snapshot commit fenced to the exact backend connection.
44    pub fn replace_order_snapshot_like_cpp_from_source(
45        &self,
46        acc_id: u64,
47        backend_orders: &[CachedOrder],
48        source: OrderRelationSourceToken,
49        source_projection_complete: bool,
50    ) -> OrderSnapshotReplacementPlan {
51        self.replace_order_snapshot_like_cpp_internal(
52            acc_id,
53            backend_orders,
54            Self::now_ms(),
55            Some(source),
56            source_projection_complete,
57        )
58    }
59
60    /// Injectable-clock variant used by local-stub boundary tests.
61    pub fn replace_order_snapshot_like_cpp_with_now(
62        &self,
63        acc_id: u64,
64        backend_orders: &[CachedOrder],
65        now_ms: u64,
66    ) -> OrderSnapshotReplacementPlan {
67        self.replace_order_snapshot_like_cpp_internal(acc_id, backend_orders, now_ms, None, true)
68    }
69
70    fn replace_order_snapshot_like_cpp_internal(
71        &self,
72        acc_id: u64,
73        backend_orders: &[CachedOrder],
74        now_ms: u64,
75        source: Option<OrderRelationSourceToken>,
76        source_projection_complete: bool,
77    ) -> OrderSnapshotReplacementPlan {
78        let relation_lock = self.order_relation_snapshot_lock(acc_id);
79        let _relation_guard = relation_lock.lock();
80        if let Some(source) = source
81            && !Self::order_relation_source_is_admitted(
82                &self.order_relation_snapshot_state(acc_id),
83                source,
84            )
85        {
86            tracing::debug!(
87                acc_id,
88                source_connection_generation = source.connection_generation,
89                "ignored stale order snapshot after relation-generation fence"
90            );
91            return OrderSnapshotReplacementPlan {
92                emit_order_indices: Vec::new(),
93            };
94        }
95        let mut entry = self.orders.entry(acc_id).or_default();
96
97        let existing_snapshot_facts = entry
98            .iter()
99            .map(Self::order_snapshot_facts)
100            .collect::<Vec<_>>();
101        let backend_snapshot_facts = backend_orders
102            .iter()
103            .map(Self::order_snapshot_facts)
104            .collect::<Vec<_>>();
105        let replacement_plan = plan_order_snapshot_replacement_like_cpp(
106            &existing_snapshot_facts,
107            &backend_snapshot_facts,
108        );
109
110        let merge_plan = {
111            let existing_facts: Vec<_> = entry.iter().map(Self::order_merge_facts).collect();
112            let backend_facts: Vec<_> =
113                backend_orders.iter().map(Self::order_merge_facts).collect();
114            plan_backend_authoritative_order_merge_like_cpp(&existing_facts, &backend_facts, now_ms)
115        };
116
117        let mut new_orders: Vec<CachedOrder> = backend_orders
118            .iter()
119            .cloned()
120            .zip(merge_plan.backend_patches)
121            .map(|(mut order, patch)| {
122                order.is_stub = patch.is_stub;
123                order.is_local_order = patch.is_local_order;
124                order.stub_inserted_at_ms = patch.stub_inserted_at_ms;
125                order.is_pending_broker_confirm = patch.is_pending_broker_confirm;
126                order
127            })
128            .collect();
129
130        for index in merge_plan.retained_existing_indices {
131            if let Some(order) = entry.get(index) {
132                new_orders.push(order.clone());
133            }
134        }
135
136        let availability = if source
137            .is_some_and(|source| source.channel == super::OrderRelationSourceChannel::Unknown)
138            || !source_projection_complete
139            || new_orders
140                .iter()
141                .any(|order| order.is_stub || order.is_pending_broker_confirm)
142        {
143            super::OrderRelationSnapshotAvailability::Stale
144        } else {
145            super::OrderRelationSnapshotAvailability::Fresh
146        };
147        *entry = new_orders;
148        drop(entry);
149        self.publish_order_relation_snapshot_state(acc_id, availability, source);
150        self.prune_order_brokers_for_acc(acc_id);
151        replacement_plan
152    }
153
154    fn order_snapshot_facts(order: &CachedOrder) -> OrderSnapshotRowFacts {
155        OrderSnapshotRowFacts {
156            order: TradeOrderFacts {
157                order_id: order.order_id,
158                backend_order_id: order.backend_order_id.clone(),
159                order_id_ex: order.order_id_ex.clone(),
160                order_status: order.order_status,
161                qty: order.qty,
162                price: order.price,
163                fill_qty: order.fill_qty,
164                aux_price: order.aux_price,
165                trail_type: order.trail_type,
166                trail_value: order.trail_value,
167                trail_spread: order.trail_spread,
168            },
169            order_version: order.order_version,
170        }
171    }
172
173    fn order_merge_facts(order: &CachedOrder) -> BackendOrderMergeOrderFacts<'_> {
174        BackendOrderMergeOrderFacts {
175            order_id: order.order_id,
176            backend_order_id: &order.backend_order_id,
177            order_id_ex: &order.order_id_ex,
178            is_stub: order.is_stub,
179            is_local_order: order.is_local_order,
180            is_local_multileg_order: order.is_local_order
181                && order.security_type
182                    == futu_proto_internal::odr_sys_cmn::SecurityType::Multileg as i32,
183            order_status: order.order_status,
184            stub_inserted_at_ms: order.stub_inserted_at_ms,
185            is_pending_broker_confirm: order.is_pending_broker_confirm,
186        }
187    }
188
189    fn now_ms() -> u64 {
190        use std::time::{SystemTime, UNIX_EPOCH};
191        SystemTime::now()
192            .duration_since(UNIX_EPOCH)
193            .map(|d| d.as_millis() as u64)
194            .unwrap_or(0)
195    }
196}