futu_backend/stock_list/
event_notice.rs1use std::collections::HashSet;
6
7use bytes::Bytes;
8use futu_core::error::{FutuError, Result};
9use futu_domain_qot_market_state::{
10 EventNoticeMarket, PULL_EVENT_NOTICE_MARKETS, SUBSCRIBE_EVENT_NOTICE_MARKETS,
11 ensure_event_notice_pull_success, ensure_event_notice_subscription_success,
12};
13use futu_domain_qot_market_state::{MarketEventUpdate, MarketStatePatch, MarketTradeDateSnapshot};
14use futures::future::join_all;
15use prost::Message;
16
17use crate::command_runtime::{MarketEventOperation, execute_market_event};
18use crate::conn::BackendConn;
19use crate::proto_internal::{ft_cmd6822, ft_cmd6823};
20
21use super::market::market_state_priority;
22use super::market::{quote_market_from_market_id, quote_mkt_to_nn_mkt_type};
23use super::{QuoteMktType, make_quote_reserved, market_id_matches};
24
25#[derive(Debug, Clone)]
26pub struct MarketStatus {
27 pub market_id: u32,
28 pub status: u32,
29 pub commit_day_key: Option<u32>,
30 pub status_text: String,
31 pub status_text_tc: String,
32}
33
34#[must_use]
42pub fn market_status_snapshot_updates(statuses: &[MarketStatus]) -> Vec<MarketEventUpdate> {
43 let mut updates = Vec::with_capacity(statuses.len() * 2);
44 for status in statuses {
45 let quote_market_type = quote_market_from_market_id(status.market_id)
46 .map(quote_mkt_to_nn_mkt_type)
47 .unwrap_or(0);
50 updates.push(MarketEventUpdate::MarketState(MarketStatePatch {
51 market_id: status.market_id,
52 quote_market_type,
53 status: status.status,
54 status_text_sc: Some(status.status_text.clone()),
55 status_text_tc: Some(status.status_text_tc.clone()),
56 }));
57 if let Some(trade_date) = status.commit_day_key {
58 updates.push(MarketEventUpdate::MarketTradeDate(
59 MarketTradeDateSnapshot {
60 market_id: status.market_id,
61 quote_market_type,
62 trade_date,
63 },
64 ));
65 }
66 }
67 updates
68}
69
70pub async fn register_markets(backend: &BackendConn) -> Result<()> {
71 for &market in SUBSCRIBE_EVENT_NOTICE_MARKETS {
72 let quote_market = quote_market_for_event_notice(market);
73 if let Err(error) = register_market(backend, quote_market).await {
74 tracing::warn!(market = ?market, error = %error, "CMD6822 registration failed");
77 }
78 }
79 Ok(())
80}
81
82pub(crate) async fn register_market(backend: &BackendConn, market: QuoteMktType) -> Result<()> {
83 let request = ft_cmd6822::RegisterReq { reserved: 0 };
86 let response = execute_market_event(
87 backend,
88 MarketEventOperation::Subscribe,
89 Bytes::from(request.encode_to_vec()),
90 make_quote_reserved(market, 0),
91 )
92 .await?;
93 let parsed = ft_cmd6822::RegisterRes::decode(response.body.as_ref())?;
94 ensure_event_notice_subscription_success(parsed.res).map_err(|reject| {
95 FutuError::Codec(format!(
96 "CMD6822 registration rejected with result {}",
97 reject.result_code()
98 ))
99 })?;
100 tracing::debug!(market = ?market, "CMD6822 registered");
101 Ok(())
102}
103
104pub async fn pull_market_status(backend: &BackendConn) -> Result<Vec<MarketStatus>> {
106 pull_all_market_status(backend).await
107}
108
109pub async fn pull_all_market_status(backend: &BackendConn) -> Result<Vec<MarketStatus>> {
110 let pulls = PULL_EVENT_NOTICE_MARKETS.iter().map(|&market| async move {
111 let quote_market = quote_market_for_event_notice(market);
112 (market, pull_market_status_for(backend, quote_market).await)
113 });
114 let results = join_all(pulls).await;
115
116 let mut seen_ids = HashSet::new();
117 let mut statuses = Vec::new();
118 for (market, result) in results {
119 match result {
120 Ok(rows) => {
121 statuses.extend(
122 rows.into_iter()
123 .filter(|row| seen_ids.insert(row.market_id)),
124 );
125 }
126 Err(error) => {
127 tracing::warn!(
130 market = ?market,
131 error = %error,
132 "CMD6823 market pull failed; keeping other market results"
133 );
134 }
135 }
136 }
137
138 tracing::debug!(
139 total_entries = statuses.len(),
140 market_types_pulled = PULL_EVENT_NOTICE_MARKETS.len(),
141 "CMD6823 multi-market snapshot fetched"
142 );
143 Ok(statuses)
144}
145
146pub(crate) async fn pull_market_status_for(
147 backend: &BackendConn,
148 market: QuoteMktType,
149) -> Result<Vec<MarketStatus>> {
150 let request = ft_cmd6823::MarketStatusReq { reserved: 0 };
153 let response = execute_market_event(
154 backend,
155 MarketEventOperation::Pull,
156 Bytes::from(request.encode_to_vec()),
157 make_quote_reserved(market, 0),
158 )
159 .await?;
160 let parsed = decode_market_status_rsp(response.body.as_ref())?;
161 ensure_event_notice_pull_success(parsed.res_status.len())
162 .map_err(|_| FutuError::Codec("CMD6823 response has no market status rows".to_string()))?;
163
164 Ok(parsed
165 .res_status
166 .into_iter()
167 .map(|status| MarketStatus {
168 market_id: status.id,
169 status: status.status.unwrap_or(0),
172 commit_day_key: status.commit_day_key,
173 status_text: status.status_text_sc.unwrap_or_default(),
174 status_text_tc: status.status_text_tc.unwrap_or_default(),
175 })
176 .collect())
177}
178
179pub fn pick_market_state(all: &[MarketStatus], market: QuoteMktType) -> Option<u32> {
180 let mut best = None;
181 let mut best_priority = 0;
182
183 for status in all
184 .iter()
185 .filter(|status| market_id_matches(market, status.market_id))
186 {
187 let priority = market_state_priority(status.status);
188 if priority > best_priority {
189 best = Some(status.status);
190 best_priority = priority;
191 } else if best.is_none() {
192 best = Some(status.status);
193 }
194 }
195 best
196}
197
198pub async fn pull_single_market_status(
199 backend: &BackendConn,
200 market: QuoteMktType,
201) -> Result<Vec<MarketStatus>> {
202 let all = pull_all_market_status(backend).await?;
203 Ok(all
204 .into_iter()
205 .filter(|status| market_id_matches(market, status.market_id))
206 .collect())
207}
208
209pub(crate) fn decode_market_status_rsp(body: &[u8]) -> Result<ft_cmd6823::MarketStatusRsp> {
210 ft_cmd6823::MarketStatusRsp::decode(body).map_err(FutuError::Proto)
211}
212
213pub(super) fn quote_market_for_event_notice(market: EventNoticeMarket) -> QuoteMktType {
214 match market {
215 EventNoticeMarket::HK => QuoteMktType::HK,
216 EventNoticeMarket::HKFuture => QuoteMktType::HKFuture,
217 EventNoticeMarket::HKFuture2 => QuoteMktType::HKFuture2,
218 EventNoticeMarket::HKOption => QuoteMktType::HKOption,
219 EventNoticeMarket::US => QuoteMktType::US,
220 EventNoticeMarket::USOption => QuoteMktType::USOption,
221 EventNoticeMarket::SH => QuoteMktType::SH,
222 EventNoticeMarket::SHKcb => QuoteMktType::SHKC,
223 EventNoticeMarket::SZ => QuoteMktType::SZ,
224 EventNoticeMarket::USFuture => QuoteMktType::USFuture,
225 EventNoticeMarket::SGFuture => QuoteMktType::SGFuture,
226 EventNoticeMarket::JPFuture => QuoteMktType::JPFuture,
227 EventNoticeMarket::JPSecurity => QuoteMktType::JPSecurity,
228 EventNoticeMarket::DigitalCcy => QuoteMktType::DigitalCcy,
229 EventNoticeMarket::SGSecurity => QuoteMktType::SGSecurity,
230 EventNoticeMarket::MYSecurity => QuoteMktType::MYSecurity,
231 EventNoticeMarket::EventContract => QuoteMktType::EventContract,
232 }
233}