Skip to main content

longbridge/market/
context.rs

1use std::sync::Arc;
2
3use longbridge_httpcli::{DcRegion, HttpClient, Json, Method};
4use serde::{Deserialize, Serialize, de::DeserializeOwned};
5use tracing::{Subscriber, dispatcher, instrument::WithSubscriber};
6
7use crate::{Config, Result, market::types::*};
8
9/// Convert a Unix-seconds value (integer or string) to RFC 3339.
10fn unix_secs_to_rfc3339(ts: i64) -> String {
11    time::OffsetDateTime::from_unix_timestamp(ts)
12        .map(|dt| {
13            use time::format_description::well_known::Rfc3339;
14            dt.format(&Rfc3339).unwrap_or_default()
15        })
16        .unwrap_or_else(|_| ts.to_string())
17}
18
19/// Convert a Unix-seconds string to RFC 3339.
20fn unix_secs_str_to_rfc3339(s: &str) -> String {
21    s.parse::<i64>()
22        .map(unix_secs_to_rfc3339)
23        .unwrap_or_else(|_| s.to_string())
24}
25
26struct InnerMarketContext {
27    http_cli: HttpClient,
28    log_subscriber: Arc<dyn Subscriber + Send + Sync>,
29}
30
31impl Drop for InnerMarketContext {
32    fn drop(&mut self) {
33        dispatcher::with_default(&self.log_subscriber.clone().into(), || {
34            tracing::info!("market context dropped");
35        });
36    }
37}
38
39/// Market data context — broker holdings, A/H premium, trade statistics,
40/// market anomalies, index constituents and more.
41#[derive(Clone)]
42pub struct MarketContext(Arc<InnerMarketContext>);
43
44impl MarketContext {
45    /// Create a [`MarketContext`]
46    pub fn new(config: Arc<Config>) -> Self {
47        let log_subscriber = config.create_log_subscriber("market");
48        dispatcher::with_default(&log_subscriber.clone().into(), || {
49            tracing::info!(language = ?config.language, "creating market context");
50        });
51        let ctx = Self(Arc::new(InnerMarketContext {
52            http_cli: config.create_http_client(),
53            log_subscriber,
54        }));
55        dispatcher::with_default(&ctx.0.log_subscriber.clone().into(), || {
56            tracing::info!("market context created");
57        });
58        ctx
59    }
60
61    /// Returns the log subscriber
62    #[inline]
63    pub fn log_subscriber(&self) -> Arc<dyn Subscriber + Send + Sync> {
64        self.0.log_subscriber.clone()
65    }
66
67    async fn get<R, Q>(&self, path: &'static str, query: Q) -> Result<R>
68    where
69        R: DeserializeOwned + Send + Sync + 'static,
70        Q: Serialize + Send + Sync,
71    {
72        Ok(self
73            .0
74            .http_cli
75            .request(Method::GET, path)
76            .query_params(query)
77            .response::<Json<R>>()
78            .send()
79            .with_subscriber(self.0.log_subscriber.clone())
80            .await?
81            .0)
82    }
83
84    /// Like [`get`](Self::get), but restricted to a single data center. Used by
85    /// region-limited endpoints (e.g. AP-only broker holdings).
86    async fn get_dc<R, Q>(&self, path: &'static str, query: Q, dc_restrict: DcRegion) -> Result<R>
87    where
88        R: DeserializeOwned + Send + Sync + 'static,
89        Q: Serialize + Send + Sync,
90    {
91        Ok(self
92            .0
93            .http_cli
94            .request(Method::GET, path)
95            .dc_restrict(dc_restrict)
96            .query_params(query)
97            .response::<Json<R>>()
98            .send()
99            .with_subscriber(self.0.log_subscriber.clone())
100            .await?
101            .0)
102    }
103
104    async fn post<R, B>(&self, path: &'static str, body: B) -> Result<R>
105    where
106        R: DeserializeOwned + Send + Sync + 'static,
107        B: std::fmt::Debug + Serialize + Send + Sync + 'static,
108    {
109        Ok(self
110            .0
111            .http_cli
112            .request(Method::POST, path)
113            .body(Json(body))
114            .response::<Json<R>>()
115            .send()
116            .with_subscriber(self.0.log_subscriber.clone())
117            .await?
118            .0)
119    }
120
121    // ── market_status ─────────────────────────────────────────────
122
123    /// Get current trading status for all markets.
124    ///
125    /// Path: `GET /v1/quote/market-status`
126    pub async fn market_status(&self) -> Result<MarketStatusResponse> {
127        #[derive(Serialize)]
128        struct Empty {}
129        self.get("/v1/quote/market-status", Empty {}).await
130    }
131
132    // ── broker_holding ────────────────────────────────────────────
133
134    /// Get top broker holdings (buy/sell leaders) for a security.
135    ///
136    /// Path: `GET /v1/quote/broker-holding`
137    pub async fn broker_holding(
138        &self,
139        symbol: impl Into<String>,
140        period: BrokerHoldingPeriod,
141    ) -> Result<BrokerHoldingTop> {
142        let period_str = match period {
143            BrokerHoldingPeriod::Rct1 => "rct_1",
144            BrokerHoldingPeriod::Rct5 => "rct_5",
145            BrokerHoldingPeriod::Rct20 => "rct_20",
146            BrokerHoldingPeriod::Rct60 => "rct_60",
147        };
148        #[derive(Serialize)]
149        struct Query {
150            symbol: String,
151            #[serde(rename = "type")]
152            period: &'static str,
153        }
154        self.get_dc(
155            "/v1/quote/broker-holding",
156            Query {
157                symbol: symbol.into(),
158                period: period_str,
159            },
160            DcRegion::Ap,
161        )
162        .await
163    }
164
165    /// Get full broker holding details for a security.
166    ///
167    /// Path: `GET /v1/quote/broker-holding/detail`
168    pub async fn broker_holding_detail(
169        &self,
170        symbol: impl Into<String>,
171    ) -> Result<BrokerHoldingDetail> {
172        #[derive(Serialize)]
173        struct Query {
174            symbol: String,
175        }
176        self.get_dc(
177            "/v1/quote/broker-holding/detail",
178            Query {
179                symbol: symbol.into(),
180            },
181            DcRegion::Ap,
182        )
183        .await
184    }
185
186    /// Get daily holding history for a specific broker.
187    ///
188    /// Path: `GET /v1/quote/broker-holding/daily`
189    pub async fn broker_holding_daily(
190        &self,
191        symbol: impl Into<String>,
192        broker_id: impl Into<String>,
193    ) -> Result<BrokerHoldingDailyHistory> {
194        #[derive(Serialize)]
195        struct Query {
196            symbol: String,
197            parti_number: String,
198        }
199        self.get_dc(
200            "/v1/quote/broker-holding/daily",
201            Query {
202                symbol: symbol.into(),
203                parti_number: broker_id.into(),
204            },
205            DcRegion::Ap,
206        )
207        .await
208    }
209
210    // ── ah_premium ────────────────────────────────────────────────
211
212    /// Get A/H premium K-line data for a dual-listed security.
213    ///
214    /// Path: `GET /v1/quote/ahpremium/klines`
215    pub async fn ah_premium(
216        &self,
217        symbol: impl Into<String>,
218        period: AhPremiumPeriod,
219        count: u32,
220    ) -> Result<AhPremiumKlines> {
221        #[derive(Serialize)]
222        struct Query {
223            symbol: String,
224            line_type: &'static str,
225            line_num: u32,
226        }
227        self.get(
228            "/v1/quote/ahpremium/klines",
229            Query {
230                symbol: symbol.into(),
231                line_type: period.to_line_type(),
232                line_num: count,
233            },
234        )
235        .await
236    }
237
238    /// Get A/H premium intraday data for a dual-listed security.
239    ///
240    /// Path: `GET /v1/quote/ahpremium/timeshares`
241    pub async fn ah_premium_intraday(
242        &self,
243        symbol: impl Into<String>,
244    ) -> Result<AhPremiumIntraday> {
245        #[derive(Serialize)]
246        struct Query {
247            symbol: String,
248            days: &'static str,
249        }
250        self.get(
251            "/v1/quote/ahpremium/timeshares",
252            Query {
253                symbol: symbol.into(),
254                days: "1",
255            },
256        )
257        .await
258    }
259
260    // ── trade_stats ───────────────────────────────────────────────
261
262    /// Get buy/sell/neutral trade statistics for a security.
263    ///
264    /// Path: `GET /v1/quote/trades-statistics`
265    pub async fn trade_stats(&self, symbol: impl Into<String>) -> Result<TradeStatsResponse> {
266        #[derive(Serialize)]
267        struct Query {
268            symbol: String,
269        }
270        self.get(
271            "/v1/quote/trades-statistics",
272            Query {
273                symbol: symbol.into(),
274            },
275        )
276        .await
277    }
278
279    // ── anomaly ───────────────────────────────────────────────────
280
281    /// Get market anomaly alerts (unusual price/volume events).
282    ///
283    /// Path: `GET /v1/quote/changes`
284    pub async fn anomaly(&self, market: impl Into<String>) -> Result<AnomalyResponse> {
285        #[derive(Serialize)]
286        struct Query {
287            market: String,
288            category: &'static str,
289        }
290        self.get(
291            "/v1/quote/changes",
292            Query {
293                market: market.into().to_uppercase(),
294                category: "0",
295            },
296        )
297        .await
298    }
299
300    // ── constituent ───────────────────────────────────────────────
301
302    /// Get constituent stocks for an index.
303    ///
304    /// `symbol` should be an index symbol such as `"HSI.HK"`.
305    ///
306    /// Path: `GET /v1/quote/index-constituents`
307    pub async fn constituent(&self, symbol: impl Into<String>) -> Result<IndexConstituents> {
308        #[derive(Serialize)]
309        struct Query {
310            symbol: String,
311        }
312        self.get(
313            "/v1/quote/index-constituents",
314            Query {
315                symbol: symbol.into(),
316            },
317        )
318        .await
319    }
320
321    // ── top_movers ────────────────────────────────────────────────
322
323    /// Get top movers (stocks with unusual price movements) across one or more
324    /// markets.
325    ///
326    /// Path: `POST /v1/quote/market/stock-events`
327    ///
328    /// `sort` is the sort order code (0 = ascending, 1 = descending).
329    /// `date` is an optional date filter in `"YYYY-MM-DD"` format.
330    pub async fn top_movers(
331        &self,
332        markets: Vec<String>,
333        sort: u32,
334        date: Option<String>,
335        limit: u32,
336    ) -> Result<TopMoversResponse> {
337        #[derive(Debug, Serialize)]
338        struct Body {
339            limit: u32,
340            sort: u32,
341            markets: Vec<String>,
342            #[serde(skip_serializing_if = "Option::is_none")]
343            date: Option<String>,
344        }
345        let raw: serde_json::Value = self
346            .post(
347                "/v1/quote/market/stock-events",
348                Body {
349                    limit,
350                    sort,
351                    markets,
352                    date,
353                },
354            )
355            .await?;
356
357        let events = raw["events"]
358            .as_array()
359            .cloned()
360            .unwrap_or_default()
361            .into_iter()
362            .map(|ev| {
363                let ts = if let Some(n) = ev["timestamp"].as_i64() {
364                    unix_secs_to_rfc3339(n)
365                } else if let Some(s) = ev["timestamp"].as_str() {
366                    unix_secs_str_to_rfc3339(s)
367                } else {
368                    String::new()
369                };
370                let stock_val = &ev["stock"];
371                let stock = TopMoversStock {
372                    symbol: stock_val["symbol"].as_str().unwrap_or("").to_string(),
373                    code: stock_val["code"].as_str().unwrap_or("").to_string(),
374                    name: stock_val["name"].as_str().unwrap_or("").to_string(),
375                    full_name: stock_val["full_name"].as_str().unwrap_or("").to_string(),
376                    change: stock_val["change"].as_str().unwrap_or("").to_string(),
377                    last_done: stock_val["last_done"].as_str().unwrap_or("").to_string(),
378                    market: stock_val["market"].as_str().unwrap_or("").to_string(),
379                    labels: stock_val["labels"]
380                        .as_array()
381                        .map(|arr| {
382                            arr.iter()
383                                .filter_map(|l| l.as_str().map(|s| s.to_string()))
384                                .collect()
385                        })
386                        .unwrap_or_default(),
387                    logo: stock_val["logo"].as_str().unwrap_or("").to_string(),
388                };
389                TopMoversEvent {
390                    timestamp: ts,
391                    alert_reason: ev["alert_reason"].as_str().unwrap_or("").to_string(),
392                    alert_type: ev["alert_type"].as_i64().unwrap_or(0),
393                    stock,
394                    post: ev["post"].clone(),
395                }
396            })
397            .collect();
398        let next_params = serde_json::to_string(&raw["next_params"]).unwrap_or_default();
399        Ok(TopMoversResponse {
400            events,
401            next_params,
402        })
403    }
404
405    // ── rank_categories ───────────────────────────────────────────
406
407    /// Get all available rank category keys and labels.
408    ///
409    /// Path: `GET /v1/quote/market/rank/categories`
410    pub async fn rank_categories(&self) -> Result<RankCategoriesResponse> {
411        #[derive(Serialize)]
412        struct Empty {}
413        #[derive(Deserialize)]
414        struct RawSubTag {
415            key: String,
416            name: String,
417            #[serde(default)]
418            market: String,
419        }
420        #[derive(Deserialize)]
421        struct RawTag {
422            key: String,
423            name: String,
424            #[serde(default)]
425            second_tags: Vec<RawSubTag>,
426        }
427        #[derive(Deserialize)]
428        struct RawData {
429            #[serde(default)]
430            first_tags: Vec<RawTag>,
431        }
432        let raw: RawData = self
433            .get("/v1/quote/market/rank/categories", Empty {})
434            .await?;
435        let categories = raw
436            .first_tags
437            .into_iter()
438            .map(|tag| {
439                let key = tag.key.strip_prefix("ib_").unwrap_or(&tag.key).to_string();
440                let sub_categories = tag
441                    .second_tags
442                    .into_iter()
443                    .map(|sub| RankSubCategory {
444                        key: sub.key.strip_prefix("ib_").unwrap_or(&sub.key).to_string(),
445                        name: sub.name,
446                        market: sub.market,
447                    })
448                    .collect();
449                RankCategory {
450                    key,
451                    name: tag.name,
452                    sub_categories,
453                }
454            })
455            .collect();
456        Ok(RankCategoriesResponse { categories })
457    }
458
459    // ── rank_list ─────────────────────────────────────────────────
460
461    /// Get a ranked list of securities for the given category key.
462    ///
463    /// Path: `GET /v1/quote/market/rank/list`
464    pub async fn rank_list(
465        &self,
466        key: impl Into<String>,
467        need_article: bool,
468    ) -> Result<RankListResponse> {
469        #[derive(Serialize)]
470        struct Query {
471            key: String,
472            delay_bmp: &'static str,
473            need_article: &'static str,
474        }
475        let key_str = key.into();
476        // Add "ib_" prefix if the caller passed a clean key (without it).
477        let api_key = if key_str.starts_with("ib_") {
478            key_str
479        } else {
480            format!("ib_{key_str}")
481        };
482        let raw: serde_json::Value = self
483            .get(
484                "/v1/quote/market/rank/list",
485                Query {
486                    key: api_key,
487                    delay_bmp: "false",
488                    need_article: if need_article { "true" } else { "false" },
489                },
490            )
491            .await?;
492        let bmp = raw["bmp"].as_bool().unwrap_or(false);
493        let lists = raw["lists"]
494            .as_array()
495            .cloned()
496            .unwrap_or_default()
497            .into_iter()
498            .map(|item| RankListItem {
499                symbol: item["symbol"].as_str().unwrap_or("").to_string(),
500                code: item["code"].as_str().unwrap_or("").to_string(),
501                name: item["name"].as_str().unwrap_or("").to_string(),
502                last_done: item["last_done"].as_str().unwrap_or("").to_string(),
503                chg: item["chg"].as_str().unwrap_or("").to_string(),
504                change: item["change"].as_str().unwrap_or("").to_string(),
505                inflow: item["inflow"].as_str().unwrap_or("").to_string(),
506                market_cap: item["market_cap"].as_str().unwrap_or("").to_string(),
507                industry: item["industry"].as_str().unwrap_or("").to_string(),
508                pre_post_price: item["pre_post_price"].as_str().unwrap_or("").to_string(),
509                pre_post_chg: item["pre_post_chg"].as_str().unwrap_or("").to_string(),
510                amplitude: item["amplitude"].as_str().unwrap_or("").to_string(),
511                five_day_chg: item["five_day_chg"].as_str().unwrap_or("").to_string(),
512                turnover_rate: item["turnover_rate"].as_str().unwrap_or("").to_string(),
513                volume_rate: item["volume_rate"].as_str().unwrap_or("").to_string(),
514                pb_ttm: item["pb_ttm"].as_str().unwrap_or("").to_string(),
515            })
516            .collect();
517        Ok(RankListResponse { bmp, lists })
518    }
519}