Skip to main content

futu_backend/stock_list/
event_notice.rs

1//! CMD6822/6823 market-event runtime adapter.
2//!
3//! C++ behavior owner: `NNBiz_Qot_EventNotice.cpp:38-229`.
4
5use 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/// Adapt the complete CMD6823 pull snapshot into the same domain updates used
35/// by CMD6301 pushes.
36///
37/// C++ `NNBiz_Qot_EventNotice.cpp:106-125` constructs a new market-state row
38/// from each pull item and stores `commit_day_key` alongside it. The pull path
39/// therefore replaces both localized text fields, including an explicit empty
40/// string, while CMD6301 patches preserve text fields absent on the wire.
41#[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            // C++ `NN_QuoteMktType_From_NN_QuoteMktID` returns UNKNOWN=0
48            // for an unclassified market id; the row is still persisted.
49            .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            // C++ sends all market registrations independently and treats a
75            // failed reply as non-fatal to the remaining registrations.
76            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    // Protocol constant, not a missing dynamic source. C++ writes zero in
84    // `NNBiz_Qot_EventNotice.cpp:175-181`; re-audit if that request changes.
85    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
104/// Startup compatibility entry. Current C++ pulls every event-notice market.
105pub 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                // C++ replies are independent: one malformed/failed market
128                // does not discard rows already applied by other replies.
129                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    // Protocol constant, not a missing dynamic source. C++ writes zero in
151    // `NNBiz_Qot_EventNotice.cpp:52-58`; re-audit if that request changes.
152    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            // C++ zero-initializes Ndt_Qot_MarketStateInfo before unpacking
170            // optional status/text fields (`NNBiz_Qot_EventNotice.cpp:106-113`).
171            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}