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
9fn 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
19fn 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#[derive(Clone)]
42pub struct MarketContext(Arc<InnerMarketContext>);
43
44impl MarketContext {
45 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}