futu_cache/trd_cache/
backend_merge.rs1use 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 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 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 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 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 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}